net_processing: add a global delay queue for sending txs

Without the per-peer rate limiting, nodes can act as an amplifier for
transaction spam -- receiving many transactions from one node, but
relaying each of them to over 100 other nodes. Limit the impact of this
by providing a global rate limit.

This is implemented using dual token buckets, one that consumes a
token for every transaction, and one that consumes a token for every
serialized byte. This rate limits both per-tx resource usage (eg INV
messages) and overall relay bandwidth.

Main bucket parameters:
 * Count: 14tx/s rate, 420tx (30s) capacity
 * Size: 12MB/600s rate (4-6 blocks per target block interval), 50MB capacity

The size bucket is expected to be large enough to almost never have an
impact in normal usage, even during transaction storms, and is primarily
intended to mitigate attack-like scenarios.

Outbound connections get a separate pair of buckets, with rates boosted
by a 2.5x multiplier.

This avoids the excessive memory and CPU usage due to the 100x multiplier
from the queues being per-peer.

Note that this also reduces the size of INV messages we send for general
tx relay back to a more reasonable level of under 600 txs in 99.999%
of cases.
This commit is contained in:
Anthony Towns
2026-07-12 07:53:06 +10:00
parent 7927650e56
commit df31ee57aa
2 changed files with 176 additions and 25 deletions

View File

@@ -59,6 +59,7 @@
#include <util/check.h>
#include <util/strencodings.h>
#include <util/time.h>
#include <util/tokenbucket.h>
#include <util/trace.h>
#include <validation.h>
@@ -172,6 +173,12 @@ static constexpr auto OUTBOUND_INVENTORY_BROADCAST_INTERVAL{2s};
[[maybe_unused]] static constexpr unsigned int INVENTORY_BROADCAST_PER_SECOND{14};
/** Target number of tx inventory items to send per transmission. */
[[maybe_unused]] static constexpr unsigned int INVENTORY_BROADCAST_TARGET = INVENTORY_BROADCAST_PER_SECOND * count_seconds(INBOUND_INVENTORY_BROADCAST_INTERVAL);
/** Multiplier for the inventory bucket rate for outbounds */
static constexpr double OUTBOUND_INVENTORY_BUCKET_MULTIPLIER{Ticks<SecondsDouble>(INBOUND_INVENTORY_BROADCAST_INTERVAL) / Ticks<SecondsDouble>(OUTBOUND_INVENTORY_BROADCAST_INTERVAL)};
/** Delay between checking inventory bucket and backlog */
static constexpr auto INVENTORY_BUCKET_CHECK_DELAY{100ms};
/** Empty backlog target capacity */
static constexpr size_t INVENTORY_BUCKET_BACKLOG_CAPACITY{300};
/** Average delay between feefilter broadcasts in seconds. */
static constexpr auto AVG_FEEFILTER_BROADCAST_INTERVAL{10min};
/** Maximum feefilter broadcast delay after significant change. */
@@ -494,6 +501,59 @@ struct CNodeState {
int64_t m_last_block_announcement{0};
};
struct InvToSendBucket {
const double count_floor{0};
std::vector<Wtxid> backlog;
util::TokenBucket<NodeClock> size_bucket;
util::TokenBucket<NodeClock> count_bucket;
/* Initialization rationale:
*
* Count bucket: Fills at rate*mult, total/initial capacity of 30s with mult=1
* Size bucket: Fills at 12MB every 600s, times mult so expected to be 6 times
* the rate at which blocks can confirm transactions, but at least 3 times that in
* the worst case. High limit to avoid triggering even with large spikes, but a
* modest initial value to ensure that frequent node restarts don't raise the limit
* too much.
* Count floor: In order to avoid sorting the global backlog too often, we ensure
* that we always remove at least an average INV message's number of transactions
* each time we do work. (Or 50kB if the size bucket is the limiting factor)
*/
static constexpr double SIZE_INIT{12'000'000}; // 12 MB initially
static constexpr double SIZE_CAP{50'000'000}; // 50 MB maximum
static constexpr double SIZE_REFILL{20'000}; // 20kB/s = 12MB/600s
static constexpr double INBOUND_COUNT_SECONDS{30}; // cap/initial at 30s/mult worth of txs
InvToSendBucket(unsigned int rate, double mult)
: count_floor{-1.0 * INVENTORY_BROADCAST_TARGET},
size_bucket(/*rate=*/SIZE_REFILL * mult, /*value=*/SIZE_INIT, /*cap=*/SIZE_CAP),
count_bucket(/*rate=*/rate * mult, /*value=*/rate * INBOUND_COUNT_SECONDS, /*cap=*/rate * INBOUND_COUNT_SECONDS)
{
}
bool avail() const
{
return !backlog.empty() && size_bucket.value() > 0 && count_bucket.value() > 0;
}
void increment(NodeClock::time_point now)
{
size_bucket.increment(now);
count_bucket.increment(now);
}
std::vector<Wtxid> TakeForProcessing(CTxMemPool& mempool) EXCLUSIVE_LOCKS_REQUIRED(mempool.cs);
bool decrement(double size)
{
bool size_ok = size_bucket.decrement(size, /*floor=*/-50e3);
bool count_ok = count_bucket.decrement(1, /*floor=*/count_floor);
return size_ok && count_ok;
}
};
class PeerManagerImpl final : public PeerManager
{
public:
@@ -520,9 +580,9 @@ public:
void FinalizeNode(const CNode& node) override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_headers_presync_mutex, !m_tx_download_mutex);
bool HasAllDesirableServiceFlags(ServiceFlags services) const override;
bool ProcessMessages(CNode& node, std::atomic<bool>& interrupt) override
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, !m_headers_presync_mutex, g_msgproc_mutex, !m_tx_download_mutex);
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, !m_headers_presync_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
bool SendMessages(CNode& node) override
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, g_msgproc_mutex, !m_tx_download_mutex);
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
/** Implement PeerManager */
void StartScheduledTasks(CScheduler& scheduler) override;
@@ -535,7 +595,7 @@ public:
std::vector<PrivateBroadcast::TxBroadcastInfo> GetPrivateBroadcastInfo() const override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex);
std::vector<CTransactionRef> AbortPrivateBroadcast(const uint256& id) override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex);
void SendPings() override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex);
void InitiateTxBroadcastToAll(const Txid& txid, const Wtxid& wtxid) override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex);
void InitiateTxBroadcastToAll(const Txid& txid, const Wtxid& wtxid) override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_inv_to_send_mutex);
node::TransactionError InitiateTxBroadcastPrivate(const CTransactionRef& tx) override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex);
void SetBestBlock(int height, std::chrono::seconds time) override
{
@@ -549,7 +609,7 @@ public:
private:
void ProcessMessage(Peer& peer, CNode& pfrom, const std::string& msg_type, DataStream& vRecv, NodeClock::time_point time_received,
const std::atomic<bool>& interruptMsgProc)
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, !m_headers_presync_mutex, g_msgproc_mutex, !m_tx_download_mutex);
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, !m_headers_presync_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
/** Consider evicting an outbound peer based on the amount of time they've been behind our tip */
void ConsiderEviction(CNode& pto, Peer& peer, std::chrono::seconds time_in_seconds) EXCLUSIVE_LOCKS_REQUIRED(cs_main, g_msgproc_mutex);
@@ -558,7 +618,7 @@ private:
void EvictExtraOutboundPeers(NodeClock::time_point now) EXCLUSIVE_LOCKS_REQUIRED(cs_main);
/** Retrieve unbroadcast transactions from the mempool and reattempt sending to peers */
void ReattemptInitialBroadcast(CScheduler& scheduler) EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex);
void ReattemptInitialBroadcast(CScheduler& scheduler) EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_inv_to_send_mutex);
/** Rebroadcast stale private transactions (already broadcast but not received back from the network). */
void ReattemptPrivateBroadcast(CScheduler& scheduler);
@@ -616,13 +676,13 @@ private:
/** Handle a transaction whose result was MempoolAcceptResult::ResultType::VALID.
* Updates m_txrequest, m_orphanage, and vExtraTxnForCompact. Also queues the tx for relay. */
void ProcessValidTx(NodeId nodeid, const CTransactionRef& tx, const std::list<CTransactionRef>& replaced_transactions)
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, g_msgproc_mutex, m_tx_download_mutex);
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, g_msgproc_mutex, m_tx_download_mutex, !m_inv_to_send_mutex);
/** Handle the results of package validation: calls ProcessValidTx and ProcessInvalidTx for
* individual transactions, and caches rejection for the package as a group.
*/
void ProcessPackageResult(const node::PackageToValidate& package_to_validate, const PackageMempoolAcceptResult& package_result)
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, g_msgproc_mutex, m_tx_download_mutex);
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, g_msgproc_mutex, m_tx_download_mutex, !m_inv_to_send_mutex);
/**
* Reconsider orphan transactions after a parent has been accepted to the mempool.
@@ -636,7 +696,7 @@ private:
* will be empty.
*/
bool ProcessOrphanTx(Peer& peer)
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, g_msgproc_mutex, !m_tx_download_mutex);
EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
/** Process a single headers message from a peer.
*
@@ -1101,6 +1161,13 @@ private:
/// The transactions to be broadcast privately.
PrivateBroadcast m_tx_for_private_broadcast;
mutable Mutex m_inv_to_send_mutex ACQUIRED_BEFORE(m_mempool.cs);
InvToSendBucket m_inbound_inv_bucket GUARDED_BY(m_inv_to_send_mutex);
InvToSendBucket m_outbound_inv_bucket GUARDED_BY(m_inv_to_send_mutex);
std::atomic<NodeClock::time_point> m_next_inv_bucket_check{NodeClock::time_point::min()};
void ProcessInvBacklog(NodeClock::time_point now, bool backlog_bumped=false) EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_inv_to_send_mutex);
};
const CNodeState* PeerManagerImpl::State(NodeId pnode) const
@@ -2039,7 +2106,9 @@ PeerManagerImpl::PeerManagerImpl(CConnman& connman, AddrMan& addrman,
m_mempool(pool),
m_txdownloadman(node::TxDownloadOptions{pool, m_rng, opts.deterministic_rng}),
m_warnings{warnings},
m_opts{opts}
m_opts{opts},
m_inbound_inv_bucket(/*rate=*/INVENTORY_BROADCAST_PER_SECOND, /*mult=*/1.0),
m_outbound_inv_bucket(/*rate=*/INVENTORY_BROADCAST_PER_SECOND, /*mult=*/OUTBOUND_INVENTORY_BUCKET_MULTIPLIER)
{
// While Erlay support is incomplete, it must be enabled explicitly via -txreconciliation.
// This argument can go away after Erlay support is complete.
@@ -2266,28 +2335,105 @@ void PeerManagerImpl::SendPings()
for(auto& it : m_peer_map) it.second->m_ping_queued = true;
}
void PeerManagerImpl::InitiateTxBroadcastToAll(const Txid& txid, const Wtxid& wtxid)
std::vector<Wtxid> InvToSendBucket::TakeForProcessing(CTxMemPool& mempool)
{
for (const PeerRef& peer_ref : GetAllPeers()) {
if (!peer_ref) continue;
Peer& peer{*peer_ref};
AssertLockHeld(mempool.cs);
auto tx_relay = peer.GetTxRelay();
if (!tx_relay) continue;
size_t n_to_take = static_cast<size_t>(std::max<double>(count_bucket.value() - count_floor, 0));
LOCK(tx_relay->m_tx_inventory_mutex);
// Only queue transactions for announcement once the version handshake
// is completed. The time of arrival for these transactions is
// otherwise at risk of leaking to a spy, if the spy is able to
// distinguish transactions received during the handshake from the rest
// in the announcement.
if (tx_relay->m_next_inv_send_time == 0s) continue;
std::vector<Wtxid> best;
const uint256& hash{peer.m_wtxid_relay ? wtxid.ToUint256() : txid.ToUint256()};
if (!tx_relay->m_tx_inventory_known_filter.contains(hash)) {
tx_relay->m_tx_inventory_to_send.push_back(wtxid);
auto itervec = mempool.ExtractBestByMiningScoreWithTopology(backlog, n_to_take);
bool tokens_left = true;
for (auto txiter : itervec) {
auto& wtxid = txiter->GetTx().GetWitnessHash();
if (tokens_left) {
best.push_back(wtxid);
if (!decrement(txiter->GetTx().ComputeTotalSize())) {
tokens_left = false;
}
} else {
backlog.push_back(wtxid);
}
}
// if the backlog is now empty, consider shrinking it if it's oversized
if (backlog.empty() && backlog.capacity() > INVENTORY_BUCKET_BACKLOG_CAPACITY) {
std::vector<Wtxid> dummy;
dummy.reserve(INVENTORY_BUCKET_BACKLOG_CAPACITY);
dummy.swap(backlog);
}
return best;
}
void PeerManagerImpl::ProcessInvBacklog(NodeClock::time_point now, bool backlog_bumped)
{
// Don't run the body of this function unless it's been a little
// while since the last run, or we just added a new tx to the backlog.
if (!backlog_bumped && now <= m_next_inv_bucket_check.load()) return;
m_next_inv_bucket_check = now + INVENTORY_BUCKET_CHECK_DELAY;
LOCK(m_inv_to_send_mutex);
m_inbound_inv_bucket.increment(now);
m_outbound_inv_bucket.increment(now);
// Early exit to skip pointlessly touching mempool lock
bool in_avail = m_inbound_inv_bucket.avail();
bool out_avail = m_outbound_inv_bucket.avail();
if (!in_avail && !out_avail) return;
std::vector<Wtxid> for_inbound;
std::vector<Wtxid> for_outbound;
{
LOCK(m_mempool.cs);
if (in_avail) for_inbound = m_inbound_inv_bucket.TakeForProcessing(m_mempool);
if (out_avail) for_outbound = m_outbound_inv_bucket.TakeForProcessing(m_mempool);
}
if (!for_inbound.empty() || !for_outbound.empty()) {
bool any_inbound_connected = false;
bool any_outbound_connected = false;
for (const PeerRef& peer_ref : GetAllPeers()) {
if (!peer_ref) continue;
Peer& peer{*peer_ref};
auto tx_relay = peer.GetTxRelay();
if (!tx_relay) continue;
LOCK(tx_relay->m_tx_inventory_mutex);
// Only queue transactions for announcement once the version handshake
// is completed. The time of arrival for these transactions is
// otherwise at risk of leaking to a spy, if the spy is able to
// distinguish transactions received during the handshake from the rest
// in the announcement.
if (tx_relay->m_next_inv_send_time == 0s) continue;
if (peer.m_is_inbound) {
any_inbound_connected = true;
} else {
any_outbound_connected = true;
}
for (auto& i : (peer.m_is_inbound ? for_inbound : for_outbound)) {
tx_relay->m_tx_inventory_to_send.push_back(i);
}
}
// if the node has no in/outbound connections, clear the corresponding backlog entirely
// this reduces wasted memory, and avoids having the bucket artificially empty for when
// future peers do connect.
if (!any_inbound_connected) m_inbound_inv_bucket.backlog.clear();
if (!any_outbound_connected) m_outbound_inv_bucket.backlog.clear();
}
}
void PeerManagerImpl::InitiateTxBroadcastToAll(const Txid&, const Wtxid& wtxid)
{
{
LOCK(m_inv_to_send_mutex);
m_inbound_inv_bucket.backlog.push_back(wtxid);
m_outbound_inv_bucket.backlog.push_back(wtxid);
}
ProcessInvBacklog(NodeClock::now(), /*backlog_bumped=*/true);
}
node::TransactionError PeerManagerImpl::InitiateTxBroadcastPrivate(const CTransactionRef& tx)
@@ -5855,6 +6001,8 @@ bool PeerManagerImpl::SendMessages(CNode& node)
MaybeSendSendHeaders(node, peer);
ProcessInvBacklog(now);
{
LOCK(cs_main);