From 110cb1128151a7a43bdc46a346cb230721fd7046 Mon Sep 17 00:00:00 2001 From: nullPointerEnjoyer Date: Wed, 7 Oct 2026 22:57:01 +0200 Subject: [PATCH 1/8] feat(p2p): add a per-peer fork download budget Downloading and validating the blocks behind announced headers is an expensive operation. A peer could previously make us perform an unbounded amount of such work over time by repeatedly announcing distinct header chains that don't extend our tip. Introduce a token-bucket budget (capacity + refill interval, both configurable via the p2p protocol config) that every non-empty header list will be charged against; header lists received during the initial block download are exempt so that node bootstrap can't be stalled. No wire-level changes: the budget only gates our own block requests. --- .../tests/addr_list_response_caching.rs | 2 + p2p/src/protocol.rs | 12 ++ p2p/src/sync/peer/fork_download_budget.rs | 191 ++++++++++++++++++ p2p/src/sync/peer/mod.rs | 1 + p2p/src/sync/tests/block_announcement.rs | 6 + p2p/src/sync/tests/header_list_request.rs | 2 + p2p/src/sync/tests/network_sync.rs | 8 + p2p/src/sync/tests/tx_announcement.rs | 2 + p2p/src/tests/unsupported_message.rs | 2 + 9 files changed, 226 insertions(+) create mode 100644 p2p/src/sync/peer/fork_download_budget.rs diff --git a/p2p/src/peer_manager/tests/addr_list_response_caching.rs b/p2p/src/peer_manager/tests/addr_list_response_caching.rs index a92eb6dd26..692cbfa902 100644 --- a/p2p/src/peer_manager/tests/addr_list_response_caching.rs +++ b/p2p/src/peer_manager/tests/addr_list_response_caching.rs @@ -238,6 +238,8 @@ async fn basic_test(#[case] seed: Seed) { fn make_p2p_config() -> P2pConfig { P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), // Note: with the default value we'd have to switch off the extra test checks // in address tables, because the test would take forever to complete. max_addr_list_response_address_count: 10.into(), diff --git a/p2p/src/protocol.rs b/p2p/src/protocol.rs index 34a167f055..d7d06e1d9d 100644 --- a/p2p/src/protocol.rs +++ b/p2p/src/protocol.rs @@ -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. /// @@ -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, } diff --git a/p2p/src/sync/peer/fork_download_budget.rs b/p2p/src/sync/peer/fork_download_budget.rs new file mode 100644 index 0000000000..1d497ac30f --- /dev/null +++ b/p2p/src/sync/peer/fork_download_budget.rs @@ -0,0 +1,191 @@ +// Copyright (c) 2023 RBB S.r.l +// opensource@mintlayer.org +// SPDX-License-Identifier: MIT +// Licensed under the MIT License; +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://github.com/mintlayer/mintlayer-core/blob/master/LICENSE +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! A per-peer budget limiting the amount of block downloads that announced headers can +//! trigger. +//! +//! Downloading and validating block bodies is expensive. The peer synchronization manager +//! already limits how many blocks can be requested at once, but that does not bound the +//! total amount of work a peer can cause over time: a peer that keeps announcing header +//! lists can, in principle, make us download and process an unlimited number of blocks, +//! one bounded request at a time. +//! +//! This budget bounds that cumulative work: every header list a peer sends us consumes +//! tokens from the peer's bucket; when the bucket is empty, further header lists are +//! deferred (with no error and no ban score) until tokens are refilled. The budget is not +//! spent during the initial block download, so that node bootstrap cannot be stalled. +//! Honest announcement traffic is tiny compared to the budget capacity (a few headers per +//! newly produced block), so regular block propagation is unaffected. + +use std::time::Duration; + +use common::primitives::time::Time; + +#[derive(Debug, Clone)] +pub struct ForkDownloadBudget { + /// Total number of tokens (block downloads) the bucket can hold. + capacity: u64, + /// How often the full capacity is refilled. + refill_interval: Duration, + /// Currently available tokens. + tokens: u64, + /// Time at which the last refill (partial or full) happened. + last_refill: Time, +} + +impl ForkDownloadBudget { + pub fn new(capacity: usize, refill_interval: Duration, now: Time) -> Self { + let capacity = capacity as u64; + Self { + capacity, + refill_interval, + tokens: capacity, + last_refill: now, + } + } + + fn refill(&mut self, now: Time) { + if self.tokens >= self.capacity { + self.last_refill = now; + return; + } + + let interval_secs = self.refill_interval.as_secs(); + if interval_secs == 0 { + // Degenerate configuration: refill instantly. + self.tokens = self.capacity; + self.last_refill = now; + return; + } + + let elapsed = now.saturating_sub(self.last_refill); + // Tokens gained since the last refill. Use u128 to avoid overflow for very long + // elapsed times. + let gained = (elapsed.as_secs() as u128) + .checked_mul(self.capacity as u128) + .and_then(|v| v.checked_div(interval_secs as u128)) + .unwrap_or(u128::MAX); + + if gained == 0 { + return; + } + + // Advance the refill time by the amount of time that corresponds to the granted + // tokens, so that remainders are not lost across calls. + let granted = gained.min(u64::MAX as u128) as u64; + self.tokens = self.tokens.saturating_add(granted).min(self.capacity); + let consumed_secs = (granted as u128) + .checked_mul(interval_secs as u128) + .and_then(|v| v.checked_div(self.capacity as u128)) + .unwrap_or(0); + self.last_refill = self + .last_refill + .saturating_duration_add(Duration::from_secs(consumed_secs as u64)); + } + + /// Try to consume `amount` tokens at time `now`. Returns `false` if the budget is + /// exhausted, in which case nothing is consumed. + pub fn try_take(&mut self, amount: usize, now: Time) -> bool { + self.refill(now); + + if self.tokens >= amount as u64 { + self.tokens -= amount as u64; + true + } else { + false + } + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use common::primitives::time::Time; + + use super::*; + + fn t(secs: u64) -> Time { + Time::from_secs_since_epoch(secs) + } + + #[test] + fn starts_full() { + let mut budget = ForkDownloadBudget::new(10, Duration::from_secs(100), t(0)); + assert!(budget.try_take(10, t(0))); + assert!(!budget.try_take(1, t(0))); + } + + #[test] + fn exhausts_and_refills_over_time() { + // Rate: 10 tokens per 100 seconds. + let mut budget = ForkDownloadBudget::new(10, Duration::from_secs(100), t(0)); + assert!(budget.try_take(10, t(0))); + // No time passed - nothing to refill. + assert!(!budget.try_take(1, t(0))); + // Half of the interval passed - half of the capacity refilled. + assert!(budget.try_take(5, t(50))); + assert!(!budget.try_take(1, t(50))); + // Another half of the interval passed - another half of the capacity refilled. + assert!(budget.try_take(5, t(100))); + assert!(!budget.try_take(6, t(100))); + // At t(150) the bucket holds another 5 tokens, i.e. it has been refilled at a + // constant rate all along. + assert!(budget.try_take(5, t(150))); + assert!(!budget.try_take(1, t(150))); + } + + #[test] + fn remainders_are_not_lost() { + // 4 tokens per 8 seconds => 0.5 tokens per second. Frequent small calls must + // accumulate fractional refills instead of discarding them. + let mut budget = ForkDownloadBudget::new(4, Duration::from_secs(8), t(0)); + assert!(budget.try_take(4, t(0))); + // After 1 second only 0.5 tokens have accumulated - not enough. + assert!(!budget.try_take(1, t(1))); + // After 2 seconds a whole token has accumulated. + assert!(budget.try_take(1, t(2))); + assert!(!budget.try_take(1, t(3))); + // After 4 seconds, another token has accumulated. + assert!(budget.try_take(1, t(4))); + assert!(!budget.try_take(1, t(5))); + // And so on at a constant 0.5 tokens per second. + assert!(budget.try_take(1, t(6))); + assert!(!budget.try_take(1, t(7))); + } + + #[test] + fn overdraw_is_rejected_atomically() { + let mut budget = ForkDownloadBudget::new(10, Duration::from_secs(100), t(0)); + assert!(!budget.try_take(11, t(0))); + // The rejected request must not consume anything. + assert!(budget.try_take(10, t(0))); + } + + #[test] + fn zero_interval_refills_instantly() { + let mut budget = ForkDownloadBudget::new(10, Duration::from_secs(0), t(0)); + assert!(budget.try_take(10, t(0))); + assert!(budget.try_take(10, t(0))); + } + + #[test] + fn time_going_backwards_is_safe() { + let mut budget = ForkDownloadBudget::new(10, Duration::from_secs(100), t(1000)); + assert!(budget.try_take(10, t(1000))); + // Refill time in the future: saturating arithmetic must not panic nor over-refill. + assert!(!budget.try_take(1, t(0))); + } +} diff --git a/p2p/src/sync/peer/mod.rs b/p2p/src/sync/peer/mod.rs index 9f2882dcb9..fade50b232 100644 --- a/p2p/src/sync/peer/mod.rs +++ b/p2p/src/sync/peer/mod.rs @@ -14,5 +14,6 @@ // limitations under the License. pub mod block_manager; +pub mod fork_download_budget; pub mod requested_transactions; pub mod transaction_manager; diff --git a/p2p/src/sync/tests/block_announcement.rs b/p2p/src/sync/tests/block_announcement.rs index 638ebfce67..f7421fa4f8 100644 --- a/p2p/src/sync/tests/block_announcement.rs +++ b/p2p/src/sync/tests/block_announcement.rs @@ -525,6 +525,8 @@ async fn send_headers_connected_to_previously_sent_headers(#[case] seed: Seed) { let time_getter = BasicTestTimeGetter::new(); let p2p_config = Arc::new(P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), max_request_blocks_count: 1.into(), msg_header_count_limit: Default::default(), @@ -628,6 +630,8 @@ async fn send_headers_connected_to_block_which_is_being_downloaded(#[case] seed: let time_getter = BasicTestTimeGetter::new(); let p2p_config = Arc::new(P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), max_request_blocks_count: 1.into(), msg_header_count_limit: Default::default(), @@ -728,6 +732,8 @@ async fn correct_pending_headers_update(#[case] seed: Seed) { let time_getter = BasicTestTimeGetter::new(); let p2p_config = Arc::new(P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), max_request_blocks_count: 2.into(), msg_header_count_limit: Default::default(), diff --git a/p2p/src/sync/tests/header_list_request.rs b/p2p/src/sync/tests/header_list_request.rs index 65d3b8a647..9e52fe4213 100644 --- a/p2p/src/sync/tests/header_list_request.rs +++ b/p2p/src/sync/tests/header_list_request.rs @@ -228,6 +228,8 @@ async fn locator_must_be_from_peers_known_best_block(#[case] seed: Seed) { log::debug!("common_blocks_count = {common_blocks_count}"); let p2p_config = Arc::new(test_p2p_config_with_protocol_config(ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), msg_header_count_limit: msg_header_count_limit.into(), max_request_blocks_count: Default::default(), diff --git a/p2p/src/sync/tests/network_sync.rs b/p2p/src/sync/tests/network_sync.rs index 910c0e9e85..7a025772af 100644 --- a/p2p/src/sync/tests/network_sync.rs +++ b/p2p/src/sync/tests/network_sync.rs @@ -61,6 +61,8 @@ async fn basic(#[case] seed: Seed) { let p2p_config = Arc::new(P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), msg_header_count_limit: 10.into(), max_request_blocks_count: 5.into(), @@ -303,6 +305,8 @@ async fn block_announcement_disconnected_headers(#[case] seed: Seed) { let p2p_config = Arc::new(P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), msg_header_count_limit: (MAX_REQUEST_BLOCKS_COUNT * 2).into(), max_request_blocks_count: MAX_REQUEST_BLOCKS_COUNT.into(), @@ -737,6 +741,8 @@ async fn process_block_interference2(#[case] seed: Seed) { let blocks = create_n_blocks(&mut rng, &mut tf, num_blocks); let p2p_config = Arc::new(test_p2p_config_with_protocol_config(ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), // Only 1 block in a BlockListRequest is allowed. max_request_blocks_count: 1.into(), @@ -910,6 +916,8 @@ async fn no_infinite_stalling_when_first_locator_cant_locate(#[case] seed: Seed) let time_getter = mocked_time_getter_seconds(cur_time); let p2p_config = Arc::new(test_p2p_config_with_protocol_config(ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), msg_header_count_limit: msg_header_count_limit.into(), max_request_blocks_count: Default::default(), diff --git a/p2p/src/sync/tests/tx_announcement.rs b/p2p/src/sync/tests/tx_announcement.rs index 81ef9d57cf..155e75f48d 100644 --- a/p2p/src/sync/tests/tx_announcement.rs +++ b/p2p/src/sync/tests/tx_announcement.rs @@ -239,6 +239,8 @@ async fn too_many_announcements(#[case] seed: Seed) { let p2p_config = Arc::new(P2pConfig { protocol_config: ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), max_peer_tx_announcements: 1.into(), msg_header_count_limit: Default::default(), diff --git a/p2p/src/tests/unsupported_message.rs b/p2p/src/tests/unsupported_message.rs index fdde0ae0d1..72e4924581 100644 --- a/p2p/src/tests/unsupported_message.rs +++ b/p2p/src/tests/unsupported_message.rs @@ -65,6 +65,8 @@ async fn unsupported_message_impl(seed: Seed, make_msg_too_big: bool) { let max_message_size = 1024; let max_message_size_for_peer = max_message_size * 2; let p2p_config = Arc::new(test_p2p_config_with_protocol_config(ProtocolConfig { + max_fork_downloads_per_peer: Default::default(), + fork_download_refill_interval: Default::default(), max_message_size: max_message_size.into(), msg_header_count_limit: Default::default(), From 409c1f757ef27cb119d2177b1a26f47b4b0873f8 Mon Sep 17 00:00:00 2001 From: nullPointerEnjoyer Date: Wed, 7 Oct 2026 22:57:09 +0200 Subject: [PATCH 2/8] feat(p2p): enforce the fork download budget in the block sync manager Charge every non-empty header list (including tip-anchored ones: a miner can produce an unlimited number of distinct valid children of our tip, so exempting them would leave the cost attack open) against the peer's budget before downloading the announced blocks. An exhausted budget defers the list: the peer is not punished (no ban score, no disconnect; it hasn't violated the protocol), the list is expected to arrive later via a scheduled header request retry that fires once per refill interval and is subject to the budget itself. Also stop requesting more headers from a peer that repeatedly sends already-known header lists while claiming it may have more of them, which would otherwise keep us issuing header requests forever at no cost to it. --- p2p/src/sync/peer/block_manager.rs | 98 +++++++++++++++++++++++++++++- 1 file changed, 97 insertions(+), 1 deletion(-) diff --git a/p2p/src/sync/peer/block_manager.rs b/p2p/src/sync/peer/block_manager.rs index e20d5107af..e0fa4bfb84 100644 --- a/p2p/src/sync/peer/block_manager.rs +++ b/p2p/src/sync/peer/block_manager.rs @@ -63,6 +63,8 @@ use crate::{ utils::oneshot_nofail, }; +use super::fork_download_budget::ForkDownloadBudget; + #[derive(Debug, Clone)] pub enum PeerBlockSyncManagerLocalEvent { /// Chainstate got new tip. @@ -94,6 +96,19 @@ pub struct PeerBlockSyncManager { /// 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