diff --git a/src/net_processing.cpp b/src/net_processing.cpp index b3e32499781..ddeac0da94a 100644 --- a/src/net_processing.cpp +++ b/src/net_processing.cpp @@ -59,6 +59,7 @@ #include #include #include +#include #include #include @@ -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(INBOUND_INVENTORY_BROADCAST_INTERVAL) / Ticks(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 backlog; + util::TokenBucket size_bucket; + util::TokenBucket 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 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& 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 GetPrivateBroadcastInfo() const override EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex); std::vector 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& 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& 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 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 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(std::max(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 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 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 for_inbound; + std::vector 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); diff --git a/test/functional/mempool_limit.py b/test/functional/mempool_limit.py index 4770e155de8..2e554bb80c5 100755 --- a/test/functional/mempool_limit.py +++ b/test/functional/mempool_limit.py @@ -5,6 +5,7 @@ """Test mempool limiting together/eviction with the wallet.""" from decimal import Decimal +import time from test_framework.mempool_util import ( fill_mempool, @@ -205,7 +206,9 @@ class MempoolLimitTest(BitcoinTestFramework): self.log.info('Check that mempoolminfee is minrelaytxfee') assert_equal(node.getmempoolinfo()['minrelaytxfee'], node.getmempoolinfo()["mempoolminfee"]) + node.setmocktime(int(time.time())-3600) fill_mempool(self, node) + node.setmocktime(0) # bump time forward so the rate limit buckets refresh and don't block broadcast # Deliberately try to create a tx with a fee less than the minimum mempool fee to assert that it does not get added to the mempool self.log.info('Create a mempool tx that will not pass mempoolminfee')