Skip to content
1 change: 1 addition & 0 deletions p2p/src/peer_manager/tests/addr_list_response_caching.rs
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,7 @@ fn make_p2p_config() -> P2pConfig {
msg_max_locator_count: Default::default(),
max_message_size: Default::default(),
max_peer_tx_announcements: Default::default(),
..Default::default()
},

bind_addresses: Default::default(),
Expand Down
12 changes: 12 additions & 0 deletions p2p/src/protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,12 @@ make_config_setting!(MaxMessageSize, usize, 10 * 1024 * 1024);
make_config_setting!(MaxPeerTxAnnouncements, usize, 5000);
make_config_setting!(MaxUnconnectedHeaders, usize, 10);
make_config_setting!(MaxAddrListResponseAddressCount, usize, 1000);
make_config_setting!(ForkDownloadLimit, usize, 2000);
make_config_setting!(
ForkDownloadRefillInterval,
std::time::Duration,
std::time::Duration::from_secs(600)
);

/// Protocol configuration. These values are supposed to be modified in tests only.
///
Expand Down Expand Up @@ -115,4 +121,10 @@ pub struct ProtocolConfig {
pub max_message_size: MaxMessageSize,
/// The maximum number of announcements (hashes) for which we haven't received transactions.
pub max_peer_tx_announcements: MaxPeerTxAnnouncements,
/// The maximum number of blocks that a peer can make us download for the header lists
/// it announces, before further announced header lists are deferred until the budget
/// is refilled.
pub max_fork_downloads_per_peer: ForkDownloadLimit,
/// How often the per-peer fork download budget is refilled.
pub fork_download_refill_interval: ForkDownloadRefillInterval,
}
120 changes: 119 additions & 1 deletion p2p/src/sync/peer/block_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
use std::{
collections::{BTreeSet, VecDeque},
mem,
time::Duration,
};

use itertools::Itertools;
Expand Down Expand Up @@ -63,6 +64,8 @@ use crate::{
utils::oneshot_nofail,
};

use super::fork_download_budget::ForkDownloadBudget;

#[derive(Debug, Clone)]
pub enum PeerBlockSyncManagerLocalEvent {
/// Chainstate got new tip.
Expand Down Expand Up @@ -94,8 +97,26 @@ pub struct PeerBlockSyncManager<T: NetworkingService> {
/// of headers less than the maximum. This is the signal to the peer that we have no more
/// headers, so it may not ask us for more of them in the future.
have_sent_all_headers: bool,
/// A per-peer budget limiting the total amount of fork (non-tip) block downloads that
/// this peer can cause us to perform. See `ForkDownloadBudget` for the rationale.
fork_download_budget: ForkDownloadBudget,
/// The number of consecutive header lists that contained no new headers while claiming
/// (by their size) that the peer may have more of them. Used to stop serving peers that
/// keep sending already-known header lists, which would otherwise make us issue header
/// requests forever, at no cost to them.
consecutive_known_full_header_lists: u32,
/// If set, a header list has been deferred due to the fork download budget and we
/// should ask the peer for its headers again once the given time is reached. Without
/// this, a deferred list would never be fetched unless the peer re-announces, which
/// cannot be relied upon (e.g. a quiet upstream).
fork_budget_retry_at: Option<Time>,
}

/// How often the block sync manager checks whether a fork download budget retry is due.
/// Only the due-ness is checked against the (possibly mocked) clock; the interval itself
/// is fixed so that the check cannot be postponed by other main loop activity.
const FORK_BUDGET_RETRY_POLL_INTERVAL: Duration = Duration::from_secs(1);

struct IncomingDataState {
/// A list of headers received via the `HeaderList` message that we haven't yet
/// requested the blocks for.
Expand Down Expand Up @@ -136,6 +157,11 @@ where
local_event_receiver: UnboundedReceiver<PeerBlockSyncManagerLocalEvent>,
time_getter: TimeGetter,
) -> Self {
let fork_download_budget = ForkDownloadBudget::new(
*p2p_config.protocol_config.max_fork_downloads_per_peer,
*p2p_config.protocol_config.fork_download_refill_interval,
time_getter.get_time(),
);
Self {
id: id.into(),
chain_config,
Expand All @@ -159,6 +185,9 @@ where
},
peer_activity: PeerActivity::new(),
have_sent_all_headers: false,
fork_download_budget,
consecutive_known_full_header_lists: 0,
fork_budget_retry_at: None,
}
}

Expand Down Expand Up @@ -209,6 +238,27 @@ where

_ = tokio::time::sleep(stalling_timeout),
if self.peer_activity.earliest_expected_activity_time().is_some() => {}

// An unconditional tick: it guarantees the loop wakes up at least once per
// interval even when there is no other activity, so that the fork budget
// retry deadline check below cannot be postponed indefinitely by idle
// waiting. When the loop is busy, the deadline check runs on every
// iteration anyway.
_ = tokio::time::sleep(FORK_BUDGET_RETRY_POLL_INTERVAL) => {}
}

// A header list deferred due to the fork download budget is retried once the
// (possibly mocked) clock passes the scheduled retry time. The check lives in
// the loop body rather than in a select! branch, so that a busy message flow
// cannot starve it.
let is_fork_budget_retry_due = match self.fork_budget_retry_at {
Some(retry_at) => self.time_getter.get_time() >= retry_at,
None => false,
};
if is_fork_budget_retry_due && self.common_services.has_service(Service::Blocks) {
self.fork_budget_retry_at = None;
log::debug!("Retrying a header request after a fork download budget deferral");
self.request_headers().await?;
}

self.handle_sync_status_change(&last_sync_status)?;
Expand Down Expand Up @@ -663,6 +713,53 @@ where
.expect("Headers shouldn't be empty")
.prev_block_id();

// Downloading and validating the blocks behind announced headers is expensive.
// To prevent a peer from making us perform an unbounded amount of such work over
// time (e.g. by repeatedly announcing distinct header chains that don't extend our
// tip), every header list costs tokens from the peer's download budget, and a list
// is only processed while the budget can pay for it. Everything sent during the
// initial block download is exempt, so that the budget cannot stall node bootstrap.
// Note that this must apply to tip-anchored lists as well: a malicious miner can
// produce an unlimited number of distinct valid children of our tip, so exempting
// them would leave the aforementioned attack open. Honest tip announcements are
// tiny (a few headers per new block) compared to the budget capacity, so they are
// unaffected in practice.
if !self.chainstate_handle.call(|c| Ok(c.is_initial_block_download())).await? {
let headers_count = headers.len();
if !self.fork_download_budget.try_take(headers_count, self.time_getter.get_time()) {
Comment on lines +728 to +729

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

bug · high
try_take is an all-or-nothing grant, and a header list can be larger than the budget capacity. Since ForkDownloadLimit (2000) equals the default HeaderLimit/msg_header_count_limit (2000), a peer may legitimately send a full-size (2000-header) list; if the bucket is not completely full at that moment, try_take fails, and because tokens only refill up to capacity, the same request can never succeed. The result is a permanent deferral of that list: every retry scheduled below re-requests headers, the peer re-sends the same oversized list, and it is deferred again — an infinite, self-sustaining loop that silently stalls synchronization from this peer with no ban score and no distinguishing diagnostic. Mitigate by capping the charged amount at min(headers_count, capacity) (granting partial progress per refill), or by granting/deferring in capacity-sized chunks, or at minimum ensure ForkDownloadLimit > HeaderLimit with a validation/debug assertion.

Suggestion:

Suggested change
let headers_count = headers.len();
if !self.fork_download_budget.try_take(headers_count, self.time_getter.get_time()) {
let headers_count = headers.len().min(
*self.p2p_config.protocol_config.max_fork_downloads_per_peer,
);
if !self.fork_download_budget.try_take(headers_count, self.time_getter.get_time()) {

log::info!(
"Deferring the processing of a header list of {headers_count} headers:\
the peer's fork download budget is exhausted"
);
// No error is returned and no ban score is issued: the peer hasn't violated
// the protocol, it has merely exhausted a resource quota. Also note that
// `expecting_headers_since` has already been reset at the beginning of this
// function, so this deferral cannot be mistaken for peer stalling.
//
// Schedule a header request retry for when the budget is likely to have
// recovered. This is what makes the deferral safe for liveness: without it,
// the deferred data would only be fetched if the peer re-announced it. The
// retry is a regular header request, so it's rate-limited by the refill
// interval and is subject to the budget itself.
//
// Schedule the retry no later than any retry that was already scheduled: a
// fresh deferral must not postpone a retry the budget may already have
// recovered for. This also keeps the schedule sane when the clock doesn't
// advance (as with a mocked time getter in tests).
let retry_at_candidate = self.time_getter.get_time().saturating_duration_add(
*self.p2p_config.protocol_config.fork_download_refill_interval,
);
self.fork_budget_retry_at = Some(match self.fork_budget_retry_at {
Some(existing) => existing.min(retry_at_candidate),
None => retry_at_candidate,
});
// Note: the retry is deliberately left armed here. This charge may be for a
// list different from the one the pending retry was scheduled for, and a
// deferred list can only be recovered through that retry.
return Ok(());
}
}

// Note: we require a peer to send headers starting from a block that we already have
// in our chainstate. I.e. we don't allow:
// 1) Basing new headers on a previously sent header, because this would give a malicious
Expand Down Expand Up @@ -741,11 +838,32 @@ where

if new_block_headers.is_empty() {
if peer_may_have_more_headers {
self.request_headers().await?;
// A peer that keeps sending us already-known header lists while claiming that
// it may have more of them would otherwise make us issue header requests
// forever, at no cost to itself. After a few such lists in a row, stop
// requesting until the peer sends us something actually new. This is not a
// protocol violation, so the peer isn't punished; any new headers reset the
// counter.
const MAX_CONSECUTIVE_KNOWN_FULL_HEADER_LISTS: u32 = 3;
self.consecutive_known_full_header_lists =
self.consecutive_known_full_header_lists.saturating_add(1);
if self.consecutive_known_full_header_lists
<= MAX_CONSECUTIVE_KNOWN_FULL_HEADER_LISTS
{
self.request_headers().await?;
} else {
log::info!(
"Not requesting more headers from the peer:\
it keeps sending already-known header lists"
);
}
} else {
self.consecutive_known_full_header_lists = 0;
}
return Ok(());
}

self.consecutive_known_full_header_lists = 0;
self.request_blocks(new_block_headers)
}

Expand Down
Loading
Loading