76#include <initializer_list>
89#include <unordered_set>
214 std::unique_ptr<PartiallyDownloadedBlock> partialBlock;
248 std::atomic<ServiceFlags> m_their_services{
NODE_NONE};
251 const bool m_is_inbound;
254 Mutex m_misbehavior_mutex;
256 bool m_should_discourage
GUARDED_BY(m_misbehavior_mutex){
false};
259 Mutex m_block_inv_mutex;
263 std::vector<uint256> m_blocks_for_inv_relay
GUARDED_BY(m_block_inv_mutex);
267 std::vector<uint256> m_blocks_for_headers_relay
GUARDED_BY(m_block_inv_mutex);
278 std::atomic<uint64_t> m_ping_nonce_sent{0};
282 std::atomic<bool> m_ping_queued{
false};
285 std::atomic<bool> m_wtxid_relay{
false};
297 bool m_relay_txs
GUARDED_BY(m_bloom_filter_mutex){
false};
299 std::unique_ptr<CBloomFilter> m_bloom_filter
PT_GUARDED_BY(m_bloom_filter_mutex)
GUARDED_BY(m_bloom_filter_mutex){
nullptr};
310 std::vector<Wtxid> m_tx_inventory_to_send
GUARDED_BY(m_tx_inventory_mutex);
314 bool m_send_mempool
GUARDED_BY(m_tx_inventory_mutex){
false};
317 std::chrono::microseconds m_next_inv_send_time
GUARDED_BY(m_tx_inventory_mutex){0};
320 uint64_t m_last_inv_sequence
GUARDED_BY(m_tx_inventory_mutex){1};
323 std::atomic<CAmount> m_fee_filter_received{0};
329 LOCK(m_tx_relay_mutex);
331 m_tx_relay = std::make_unique<Peer::TxRelay>();
332 return m_tx_relay.get();
337 return WITH_LOCK(m_tx_relay_mutex,
return m_tx_relay.get());
366 std::atomic_bool m_addr_relay_enabled{
false};
370 mutable Mutex m_addr_send_times_mutex;
372 std::chrono::microseconds m_next_addr_send
GUARDED_BY(m_addr_send_times_mutex){0};
374 std::chrono::microseconds m_next_local_addr_send
GUARDED_BY(m_addr_send_times_mutex){0};
377 std::atomic_bool m_wants_addrv2{
false};
386 std::atomic<uint64_t> m_addr_rate_limited{0};
388 std::atomic<uint64_t> m_addr_processed{0};
394 Mutex m_getdata_requests_mutex;
396 std::deque<CInv> m_getdata_requests
GUARDED_BY(m_getdata_requests_mutex);
402 Mutex m_headers_sync_mutex;
405 std::unique_ptr<HeadersSyncState> m_headers_sync
PT_GUARDED_BY(m_headers_sync_mutex)
GUARDED_BY(m_headers_sync_mutex) {};
408 std::atomic<bool> m_sent_sendheaders{
false};
418 std::atomic<std::chrono::seconds> m_time_offset{0
s};
422 , m_our_services{our_services}
423 , m_is_inbound{is_inbound}
427 mutable Mutex m_tx_relay_mutex;
430 std::unique_ptr<TxRelay> m_tx_relay
GUARDED_BY(m_tx_relay_mutex);
433using PeerRef = std::shared_ptr<Peer>;
445 uint256 hashLastUnknownBlock{};
451 bool fSyncStarted{
false};
453 std::chrono::microseconds m_stalling_since{0us};
454 std::list<QueuedBlock> vBlocksInFlight;
456 std::chrono::microseconds m_downloading_since{0us};
458 bool fPreferredDownload{
false};
460 bool m_requested_hb_cmpctblocks{
false};
462 bool m_provides_cmpctblocks{
false};
488 struct ChainSyncTimeoutState {
490 std::chrono::seconds m_timeout{0
s};
494 bool m_sent_getheaders{
false};
496 bool m_protect{
false};
499 ChainSyncTimeoutState m_chain_sync;
502 int64_t m_last_block_announcement{0};
505struct InvToSendBucket {
506 const double count_floor{0};
507 std::vector<Wtxid> backlog;
524 static constexpr double SIZE_INIT{12'000'000};
525 static constexpr double SIZE_CAP{50'000'000};
526 static constexpr double SIZE_REFILL{20'000};
528 static constexpr double INBOUND_COUNT_SECONDS{30};
530 InvToSendBucket(
unsigned int rate,
double mult)
532 size_bucket(SIZE_REFILL * mult, SIZE_INIT, SIZE_CAP),
533 count_bucket(rate * mult, rate * INBOUND_COUNT_SECONDS, rate * INBOUND_COUNT_SECONDS)
539 return !backlog.empty() && size_bucket.
value() > 0 && count_bucket.
value() > 0;
550 bool decrement(
double size)
552 bool size_ok = size_bucket.
decrement(size, -50e3);
553 bool count_ok = count_bucket.
decrement(1, count_floor);
554 return size_ok && count_ok;
561 .count_bucket = count_bucket.
value(),
562 .size_bucket = size_bucket.
value(),
593 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);
595 EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
610 void SetBestBlock(
int height,
std::chrono::seconds time)
override
612 m_best_height = height;
613 m_best_block_time = time;
621 const std::atomic<bool>& interruptMsgProc)
622 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);
634 void ReattemptPrivateBroadcast(
CScheduler& scheduler);
649 void Misbehaving(Peer& peer, const
std::
string& message);
660 bool via_compact_block, const
std::
string& message = "")
669 bool MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer);
682 bool MaybeDisconnectForTxRelayCapacity(
CNode&
node, const
std::
string& msg_type,
697 bool first_time_failure)
722 bool ProcessOrphanTx(Peer& peer)
732 void ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
734 bool via_compact_block)
765 bool IsContinuationOfLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
779 bool TryLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
794 void HeadersDirectFetchBlocks(
CNode& pfrom, const Peer& peer, const
CBlockIndex& last_header);
796 void UpdatePeerStateForReceivedHeaders(
CNode& pfrom, const
CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
803 template <
typename... Args>
804 void MakeAndPushMessage(
CNode&
node, std::string msg_type, Args&&...
args)
const
808 template <
typename... Args>
809 [[maybe_unused]]
void MakeAndPushFeature(
CNode&
node, std::string_view feature_id, Args&&...
args)
const
812 std::vector<unsigned char> feature_data;
819 void PushNodeVersion(
CNode& pnode,
const Peer& peer);
870 std::unique_ptr<TxReconciliationTracker> m_txreconciliation;
873 std::atomic<int> m_best_height{-1};
875 std::atomic<std::chrono::seconds> m_best_block_time{0
s};
883 const Options m_opts;
885 bool RejectIncomingTxs(
const CNode& peer)
const;
893 mutable Mutex m_peer_mutex;
900 std::map<NodeId, PeerRef> m_peer_map
GUARDED_BY(m_peer_mutex);
910 uint32_t GetFetchFlags(
const Peer& peer)
const;
912 std::map<uint64_t, std::chrono::microseconds> m_next_inv_to_inbounds_per_network_key
GUARDED_BY(g_msgproc_mutex);
929 std::atomic<int> m_wtxid_relay_peers{0};
947 std::chrono::microseconds NextInvToInbounds(std::chrono::microseconds now,
948 std::chrono::seconds average_interval,
953 Mutex m_most_recent_block_mutex;
954 std::shared_ptr<const CBlock> m_most_recent_block
GUARDED_BY(m_most_recent_block_mutex);
955 std::shared_ptr<const CBlockHeaderAndShortTxIDs> m_most_recent_compact_block
GUARDED_BY(m_most_recent_block_mutex);
957 std::unique_ptr<const std::map<GenTxid, CTransactionRef>> m_most_recent_block_txs
GUARDED_BY(m_most_recent_block_mutex);
961 Mutex m_headers_presync_mutex;
969 using HeadersPresyncStats = std::pair<arith_uint256, std::optional<std::pair<int64_t, uint32_t>>>;
971 std::map<NodeId, HeadersPresyncStats> m_headers_presync_stats
GUARDED_BY(m_headers_presync_mutex) {};
975 std::atomic_bool m_headers_presync_should_signal{
false};
1045 std::atomic<
std::chrono::seconds> m_last_tip_update{0
s};
1051 void ProcessGetData(
CNode& pfrom, Peer& peer,
const std::atomic<bool>& interruptMsgProc)
1056 void ProcessBlock(
CNode&
node,
const std::shared_ptr<const CBlock>& block,
bool force_processing,
bool min_pow_checked);
1089 std::vector<std::pair<Wtxid, CTransactionRef>> vExtraTxnForCompact
GUARDED_BY(g_msgproc_mutex);
1091 size_t vExtraTxnForCompactIt
GUARDED_BY(g_msgproc_mutex) = 0;
1103 int64_t ApproximateBestBlockDepth() const;
1113 void ProcessGetBlockData(
CNode& pfrom, Peer& peer, const
CInv& inv)
1131 bool PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
1133 const
uint256& stop_hash, uint32_t max_height_diff,
1180 void ProcessAddrs(
std::string_view msg_type,
CNode& pfrom, Peer& peer,
std::vector<
CAddress>&& vAddr, const
std::atomic<
bool>& interruptMsgProc)
1186 void LogBlockHeader(const
CBlockIndex& index, const
CNode& peer,
bool via_compact_block);
1192 InvToSendBucket m_inbound_inv_bucket
GUARDED_BY(m_inv_to_send_mutex);
1193 InvToSendBucket m_outbound_inv_bucket
GUARDED_BY(m_inv_to_send_mutex);
1194 std::atomic<
NodeClock::time_point> m_next_inv_bucket_check{NodeClock::time_point::min()};
1195 std::optional<NodeClock::time_point> m_next_inv_bucket_heartbeat
GUARDED_BY(m_inv_to_send_mutex);
1200const CNodeState* PeerManagerImpl::
State(
NodeId pnode)
const
1202 std::map<NodeId, CNodeState>::const_iterator it = m_node_states.find(pnode);
1203 if (it == m_node_states.end())
1210 return const_cast<CNodeState*
>(std::as_const(*this).State(pnode));
1218static bool IsAddrCompatible(
const Peer& peer,
const CAddress& addr)
1223void PeerManagerImpl::AddAddressKnown(Peer& peer,
const CAddress& addr)
1225 assert(peer.m_addr_known);
1226 peer.m_addr_known->insert(addr.
GetKey());
1229void PeerManagerImpl::PushAddress(Peer& peer,
const CAddress& addr)
1234 assert(peer.m_addr_known);
1235 if (addr.
IsValid() && !peer.m_addr_known->contains(addr.
GetKey()) && IsAddrCompatible(peer, addr)) {
1237 peer.m_addrs_to_send[m_rng.randrange(peer.m_addrs_to_send.size())] = addr;
1239 peer.m_addrs_to_send.push_back(addr);
1244static void AddKnownTx(Peer& peer,
const uint256& hash)
1246 auto tx_relay = peer.GetTxRelay();
1247 if (!tx_relay)
return;
1249 LOCK(tx_relay->m_tx_inventory_mutex);
1250 tx_relay->m_tx_inventory_known_filter.insert(hash);
1254static bool CanServeBlocks(
const Peer& peer)
1261static bool IsLimitedPeer(
const Peer& peer)
1268static bool CanServeWitnesses(
const Peer& peer)
1273std::chrono::microseconds PeerManagerImpl::NextInvToInbounds(std::chrono::microseconds now,
1274 std::chrono::seconds average_interval,
1275 uint64_t network_key)
1277 auto [it, inserted] = m_next_inv_to_inbounds_per_network_key.try_emplace(network_key, 0us);
1278 auto& timer{it->second};
1280 timer = now + m_rng.rand_exp_duration(average_interval);
1285bool PeerManagerImpl::IsBlockRequested(
const uint256& hash)
1287 return mapBlocksInFlight.contains(hash);
1290bool PeerManagerImpl::IsBlockRequestedFromOutbound(
const uint256& hash)
1292 for (
auto range = mapBlocksInFlight.equal_range(hash); range.first != range.second; range.first++) {
1293 auto [nodeid, block_it] = range.first->second;
1294 PeerRef peer{GetPeerRef(nodeid)};
1295 if (peer && !peer->m_is_inbound)
return true;
1301void PeerManagerImpl::RemoveBlockRequest(
const uint256& hash, std::optional<NodeId> from_peer)
1303 auto range = mapBlocksInFlight.equal_range(hash);
1304 if (range.first == range.second) {
1312 while (range.first != range.second) {
1313 const auto& [node_id, list_it]{range.first->second};
1315 if (from_peer && *from_peer != node_id) {
1322 if (state.vBlocksInFlight.begin() == list_it) {
1324 state.m_downloading_since = std::max(state.m_downloading_since, GetTime<std::chrono::microseconds>());
1326 state.vBlocksInFlight.erase(list_it);
1328 if (state.vBlocksInFlight.empty()) {
1330 m_peers_downloading_from--;
1332 state.m_stalling_since = 0us;
1334 range.first = mapBlocksInFlight.erase(range.first);
1338bool PeerManagerImpl::BlockRequested(
NodeId nodeid,
const CBlockIndex& block, std::list<QueuedBlock>::iterator** pit)
1342 CNodeState *state =
State(nodeid);
1343 assert(state !=
nullptr);
1348 for (
auto range = mapBlocksInFlight.equal_range(hash); range.first != range.second; range.first++) {
1349 if (range.first->second.first == nodeid) {
1351 *pit = &range.first->second.second;
1358 RemoveBlockRequest(hash, nodeid);
1360 std::list<QueuedBlock>::iterator it = state->vBlocksInFlight.insert(state->vBlocksInFlight.end(),
1361 {&block, std::unique_ptr<PartiallyDownloadedBlock>(pit ? new PartiallyDownloadedBlock(&m_mempool) : nullptr)});
1362 if (state->vBlocksInFlight.size() == 1) {
1364 state->m_downloading_since = GetTime<std::chrono::microseconds>();
1365 m_peers_downloading_from++;
1367 auto itInFlight = mapBlocksInFlight.insert(std::make_pair(hash, std::make_pair(nodeid, it)));
1369 *pit = &itInFlight->second.second;
1374void PeerManagerImpl::MaybeSetPeerAsAnnouncingHeaderAndIDs(
NodeId nodeid)
1381 if (m_opts.ignore_incoming_txs)
return;
1383 CNodeState* nodestate =
State(nodeid);
1384 PeerRef peer{GetPeerRef(nodeid)};
1385 if (!nodestate || !nodestate->m_provides_cmpctblocks) {
1390 int num_outbound_hb_peers = 0;
1391 for (std::list<NodeId>::iterator it = lNodesAnnouncingHeaderAndIDs.begin(); it != lNodesAnnouncingHeaderAndIDs.end(); it++) {
1392 if (*it == nodeid) {
1393 lNodesAnnouncingHeaderAndIDs.erase(it);
1394 lNodesAnnouncingHeaderAndIDs.push_back(nodeid);
1397 PeerRef peer_ref{GetPeerRef(*it)};
1398 if (peer_ref && !peer_ref->m_is_inbound) ++num_outbound_hb_peers;
1400 if (peer && peer->m_is_inbound) {
1403 if (lNodesAnnouncingHeaderAndIDs.size() >= 3 && num_outbound_hb_peers == 1) {
1404 PeerRef remove_peer{GetPeerRef(lNodesAnnouncingHeaderAndIDs.front())};
1405 if (remove_peer && !remove_peer->m_is_inbound) {
1408 std::swap(lNodesAnnouncingHeaderAndIDs.front(), *std::next(lNodesAnnouncingHeaderAndIDs.begin()));
1417 lNodesAnnouncingHeaderAndIDs.push_back(pfrom->
GetId());
1420 if (nodeid_was_appended && lNodesAnnouncingHeaderAndIDs.size() > 3) {
1423 m_connman.
ForNode(lNodesAnnouncingHeaderAndIDs.front(), [
this](
CNode* pnodeStop) {
1426 pnodeStop->m_bip152_highbandwidth_to =
false;
1429 lNodesAnnouncingHeaderAndIDs.pop_front();
1433bool PeerManagerImpl::TipMayBeStale()
1437 if (m_last_tip_update.load() == 0
s) {
1438 m_last_tip_update = GetTime<std::chrono::seconds>();
1440 return m_last_tip_update.load() < GetTime<std::chrono::seconds>() - std::chrono::seconds{consensusParams.
nPowTargetSpacing * 3} && mapBlocksInFlight.empty();
1443int64_t PeerManagerImpl::ApproximateBestBlockDepth()
const
1448bool PeerManagerImpl::CanDirectFetch()
1455 if (state->pindexBestKnownBlock && pindex == state->pindexBestKnownBlock->GetAncestor(pindex->nHeight))
1457 if (state->pindexBestHeaderSent && pindex == state->pindexBestHeaderSent->GetAncestor(pindex->nHeight))
1462void PeerManagerImpl::ProcessBlockAvailability(
NodeId nodeid) {
1463 CNodeState *state =
State(nodeid);
1464 assert(state !=
nullptr);
1466 if (!state->hashLastUnknownBlock.IsNull()) {
1469 if (state->pindexBestKnownBlock ==
nullptr || pindex->
nChainWork >= state->pindexBestKnownBlock->nChainWork) {
1470 state->pindexBestKnownBlock = pindex;
1472 state->hashLastUnknownBlock.SetNull();
1477void PeerManagerImpl::UpdateBlockAvailability(
NodeId nodeid,
const uint256 &hash) {
1478 CNodeState *state =
State(nodeid);
1479 assert(state !=
nullptr);
1481 ProcessBlockAvailability(nodeid);
1486 if (state->pindexBestKnownBlock ==
nullptr || pindex->
nChainWork >= state->pindexBestKnownBlock->nChainWork) {
1487 state->pindexBestKnownBlock = pindex;
1491 state->hashLastUnknownBlock = hash;
1496void PeerManagerImpl::FindNextBlocksToDownload(
const Peer& peer,
unsigned int count, std::vector<const CBlockIndex*>& vBlocks,
NodeId& nodeStaller)
1501 vBlocks.reserve(vBlocks.size() +
count);
1502 CNodeState *state =
State(peer.m_id);
1503 assert(state !=
nullptr);
1506 ProcessBlockAvailability(peer.m_id);
1508 if (state->pindexBestKnownBlock ==
nullptr || state->pindexBestKnownBlock->nChainWork < m_chainman.
ActiveChain().
Tip()->
nChainWork || state->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
1519 state->pindexBestKnownBlock->GetAncestor(snap_base->nHeight) != snap_base) {
1520 LogDebug(
BCLog::NET,
"Not downloading blocks from peer=%d, which doesn't have the snapshot block in its best chain.\n", peer.m_id);
1529 if (state->pindexLastCommonBlock ==
nullptr ||
1530 fork_point->nChainWork > state->pindexLastCommonBlock->nChainWork ||
1531 state->pindexBestKnownBlock->GetAncestor(state->pindexLastCommonBlock->nHeight) != state->pindexLastCommonBlock) {
1532 state->pindexLastCommonBlock = fork_point;
1534 if (state->pindexLastCommonBlock == state->pindexBestKnownBlock)
1537 const CBlockIndex *pindexWalk = state->pindexLastCommonBlock;
1543 FindNextBlocks(vBlocks, peer, state, pindexWalk,
count, nWindowEnd, &m_chainman.
ActiveChain(), &nodeStaller);
1546void PeerManagerImpl::TryDownloadingHistoricalBlocks(
const Peer& peer,
unsigned int count, std::vector<const CBlockIndex*>& vBlocks,
const CBlockIndex *from_tip,
const CBlockIndex* target_block)
1551 if (vBlocks.size() >=
count) {
1555 vBlocks.reserve(
count);
1558 if (state->pindexBestKnownBlock ==
nullptr || state->pindexBestKnownBlock->GetAncestor(target_block->
nHeight) != target_block) {
1575void PeerManagerImpl::FindNextBlocks(std::vector<const CBlockIndex*>& vBlocks,
const Peer& peer, CNodeState *state,
const CBlockIndex *pindexWalk,
unsigned int count,
int nWindowEnd,
const CChain* activeChain,
NodeId* nodeStaller)
1577 std::vector<const CBlockIndex*> vToFetch;
1578 int nMaxHeight = std::min<int>(state->pindexBestKnownBlock->nHeight, nWindowEnd + 1);
1579 bool is_limited_peer = IsLimitedPeer(peer);
1581 while (pindexWalk->
nHeight < nMaxHeight) {
1585 int nToFetch = std::min(nMaxHeight - pindexWalk->
nHeight, std::max<int>(
count - vBlocks.size(), 128));
1586 vToFetch.resize(nToFetch);
1587 pindexWalk = state->pindexBestKnownBlock->
GetAncestor(pindexWalk->
nHeight + nToFetch);
1588 vToFetch[nToFetch - 1] = pindexWalk;
1589 for (
unsigned int i = nToFetch - 1; i > 0; i--) {
1590 vToFetch[i - 1] = vToFetch[i]->
pprev;
1610 state->pindexLastCommonBlock = pindex;
1617 if (waitingfor == -1) {
1619 waitingfor = mapBlocksInFlight.lower_bound(pindex->
GetBlockHash())->second.first;
1625 if (pindex->
nHeight > nWindowEnd) {
1627 if (vBlocks.size() == 0 && waitingfor != peer.m_id) {
1629 if (nodeStaller) *nodeStaller = waitingfor;
1639 vBlocks.push_back(pindex);
1640 if (vBlocks.size() ==
count) {
1649void PeerManagerImpl::PushNodeVersion(
CNode& pnode,
const Peer& peer)
1651 uint64_t my_services;
1653 uint64_t your_services;
1655 std::string my_user_agent;
1663 my_user_agent =
"/pynode:0.0.1/";
1665 my_tx_relay =
false;
1668 my_services = peer.m_our_services;
1669 my_time = TicksSinceEpoch<std::chrono::seconds>(
NodeClock::now());
1673 my_height = m_best_height;
1674 my_tx_relay = !RejectIncomingTxs(pnode);
1693 BCLog::NET,
"send version message: version=%d, blocks=%d%s, txrelay=%d, peer=%d\n",
1696 my_tx_relay, pnode.
GetId());
1699void PeerManagerImpl::UpdateLastBlockAnnounceTime(
NodeId node, int64_t time_in_seconds)
1703 if (state) state->m_last_block_announcement = time_in_seconds;
1711 m_node_states.try_emplace(m_node_states.end(), nodeid);
1713 WITH_LOCK(m_tx_download_mutex, m_txdownloadman.CheckIsEmpty(nodeid));
1719 PeerRef peer = std::make_shared<Peer>(nodeid, our_services,
node.IsInboundConn());
1722 m_peer_map.emplace_hint(m_peer_map.end(), nodeid, peer);
1726void PeerManagerImpl::ReattemptInitialBroadcast(
CScheduler& scheduler)
1730 for (
const auto& txid : unbroadcast_txids) {
1733 if (tx !=
nullptr) {
1734 InitiateTxBroadcastToAll(tx->GetWitnessHash());
1743 scheduler.
scheduleFromNow([&] { ReattemptInitialBroadcast(scheduler); }, delta);
1746void PeerManagerImpl::ReattemptPrivateBroadcast(
CScheduler& scheduler)
1750 size_t num_for_rebroadcast{0};
1751 const auto stale_txs = m_tx_for_private_broadcast.GetStale();
1752 if (!stale_txs.empty()) {
1753 for (
const auto& stale_tx : stale_txs) {
1759 "Reattempting broadcast of stale txid=%s wtxid=%s",
1760 stale_tx->GetHash().ToString(), stale_tx->GetWitnessHash().ToString());
1761 ++num_for_rebroadcast;
1764 stale_tx->GetHash().ToString(), stale_tx->GetWitnessHash().ToString(),
1765 mempool_acceptable.m_state.ToString());
1766 m_tx_for_private_broadcast.Remove(stale_tx);
1775 scheduler.
scheduleFromNow([&] { ReattemptPrivateBroadcast(scheduler); }, delta);
1778void PeerManagerImpl::FinalizeNode(
const CNode&
node)
1789 PeerRef peer = RemovePeer(nodeid);
1791 m_wtxid_relay_peers -= peer->m_wtxid_relay;
1792 assert(m_wtxid_relay_peers >= 0);
1794 CNodeState *state =
State(nodeid);
1795 assert(state !=
nullptr);
1797 if (state->fSyncStarted)
1800 for (
const QueuedBlock& entry : state->vBlocksInFlight) {
1801 auto range = mapBlocksInFlight.equal_range(entry.pindex->GetBlockHash());
1802 while (range.first != range.second) {
1803 auto [node_id, list_it] = range.first->second;
1804 if (node_id != nodeid) {
1807 range.first = mapBlocksInFlight.erase(range.first);
1812 LOCK(m_tx_download_mutex);
1813 m_txdownloadman.DisconnectedPeer(nodeid);
1815 if (m_txreconciliation) m_txreconciliation->ForgetPeer(nodeid);
1816 m_num_preferred_download_peers -= state->fPreferredDownload;
1817 m_peers_downloading_from -= (!state->vBlocksInFlight.empty());
1818 assert(m_peers_downloading_from >= 0);
1819 m_outbound_peers_with_protect_from_disconnect -= state->m_chain_sync.m_protect;
1820 assert(m_outbound_peers_with_protect_from_disconnect >= 0);
1822 m_node_states.erase(nodeid);
1824 if (m_node_states.empty()) {
1826 assert(mapBlocksInFlight.empty());
1827 assert(m_num_preferred_download_peers == 0);
1828 assert(m_peers_downloading_from == 0);
1829 assert(m_outbound_peers_with_protect_from_disconnect == 0);
1830 assert(m_wtxid_relay_peers == 0);
1831 WITH_LOCK(m_tx_download_mutex, m_txdownloadman.CheckIsEmpty());
1834 if (
node.fSuccessfullyConnected &&
1835 !
node.IsBlockOnlyConn() && !
node.IsPrivateBroadcastConn() && !
node.IsInboundConn()) {
1843 LOCK(m_headers_presync_mutex);
1844 m_headers_presync_stats.erase(nodeid);
1846 if (
node.IsPrivateBroadcastConn() &&
1847 !m_tx_for_private_broadcast.DidNodeConfirmReception(nodeid) &&
1848 m_tx_for_private_broadcast.HavePendingTransactions()) {
1855bool PeerManagerImpl::HasAllDesirableServiceFlags(
ServiceFlags services)
const
1858 return !(GetDesirableServiceFlags(services) & (~services));
1872PeerRef PeerManagerImpl::GetPeerRef(
NodeId id)
const
1875 auto it = m_peer_map.find(
id);
1876 return it != m_peer_map.end() ? it->second :
nullptr;
1879PeerRef PeerManagerImpl::RemovePeer(
NodeId id)
1883 auto it = m_peer_map.find(
id);
1884 if (it != m_peer_map.end()) {
1885 ret = std::move(it->second);
1886 m_peer_map.erase(it);
1891std::vector<PeerRef> PeerManagerImpl::GetAllPeers()
const
1893 std::vector<PeerRef> peers;
1895 peers.reserve(m_peer_map.size());
1896 for (
const auto& [
_, peer] : m_peer_map) {
1897 peers.push_back(peer);
1906 const CNodeState* state =
State(nodeid);
1907 if (state ==
nullptr)
1909 stats.
nSyncHeight = state->pindexBestKnownBlock ? state->pindexBestKnownBlock->nHeight : -1;
1910 stats.
nCommonHeight = state->pindexLastCommonBlock ? state->pindexLastCommonBlock->nHeight : -1;
1911 for (
const QueuedBlock& queue : state->vBlocksInFlight) {
1917 PeerRef peer = GetPeerRef(nodeid);
1918 if (peer ==
nullptr)
return false;
1926 NodeClock::duration ping_wait{0us};
1927 if ((0 != peer->m_ping_nonce_sent) && (peer->m_ping_start.load() >
NodeClock::epoch)) {
1931 if (
auto tx_relay = peer->GetTxRelay(); tx_relay !=
nullptr) {
1934 LOCK(tx_relay->m_tx_inventory_mutex);
1936 stats.
m_inv_to_send = tx_relay->m_tx_inventory_to_send.size();
1948 LOCK(peer->m_headers_sync_mutex);
1949 if (peer->m_headers_sync) {
1958std::vector<node::TxOrphanage::OrphanInfo> PeerManagerImpl::GetOrphanTransactions()
1960 LOCK(m_tx_download_mutex);
1961 return m_txdownloadman.GetOrphanTransactions();
1966 LOCK(m_inv_to_send_mutex);
1969 .ignores_incoming_txs = m_opts.ignore_incoming_txs,
1970 .private_broadcast = m_opts.private_broadcast,
1971 .tx_send_rate = m_opts.tx_send_rate,
1972 .inbound_bucket = m_inbound_inv_bucket.info(),
1973 .outbound_bucket = m_outbound_inv_bucket.info(),
1977std::vector<PrivateBroadcast::TxBroadcastInfo> PeerManagerImpl::GetPrivateBroadcastInfo()
const
1979 return m_tx_for_private_broadcast.GetBroadcastInfo();
1982std::vector<CTransactionRef> PeerManagerImpl::AbortPrivateBroadcast(
const uint256&
id)
1984 const auto snapshot{m_tx_for_private_broadcast.GetBroadcastInfo()};
1985 std::vector<CTransactionRef> removed_txs;
1987 size_t connections_cancelled{0};
1988 for (
const auto& tx_info : snapshot) {
1990 if (tx->GetHash().ToUint256() !=
id && tx->GetWitnessHash().ToUint256() !=
id)
continue;
1991 if (
const auto peer_acks{m_tx_for_private_broadcast.Remove(tx)}) {
1992 removed_txs.push_back(tx);
2003void PeerManagerImpl::AddToCompactExtraTransactions(
const CTransactionRef& tx)
2005 if (m_opts.max_extra_txs == 0)
return;
2006 if (vExtraTxnForCompact.size() < m_opts.max_extra_txs) {
2007 if (vExtraTxnForCompact.empty()) vExtraTxnForCompact.reserve(m_opts.max_extra_txs);
2008 vExtraTxnForCompact.emplace_back(tx->GetWitnessHash(), tx);
2010 vExtraTxnForCompact[vExtraTxnForCompactIt] = std::make_pair(tx->GetWitnessHash(), tx);
2012 vExtraTxnForCompactIt = (vExtraTxnForCompactIt + 1) % m_opts.max_extra_txs;
2015void PeerManagerImpl::Misbehaving(Peer& peer,
const std::string& message)
2017 LOCK(peer.m_misbehavior_mutex);
2019 const std::string message_prefixed = message.empty() ?
"" : (
": " + message);
2020 peer.m_should_discourage =
true;
2029 bool via_compact_block,
const std::string& message)
2031 PeerRef peer{GetPeerRef(nodeid)};
2042 if (!via_compact_block) {
2043 if (peer) Misbehaving(*peer, message);
2051 if (peer && !via_compact_block && !peer->m_is_inbound) {
2052 if (peer) Misbehaving(*peer, message);
2059 if (peer) Misbehaving(*peer, message);
2063 if (peer) Misbehaving(*peer, message);
2068 if (message !=
"") {
2073bool PeerManagerImpl::BlockRequestAllowed(
const CBlockIndex& block_index)
2093 PeerRef peer = GetPeerRef(peer_id);
2100 RemoveBlockRequest(block_index.
GetBlockHash(), std::nullopt);
2103 if (!BlockRequested(peer_id, block_index))
return util::Unexpected{
"Already requested from this peer"};
2126 return std::make_unique<PeerManagerImpl>(connman, addrman, banman, chainman, pool, warnings, opts);
2132 : m_rng{opts.deterministic_rng},
2134 m_chainparams(chainman.GetParams()),
2138 m_chainman(chainman),
2140 m_txdownloadman{
node::TxDownloadOptions{pool, opts.deterministic_rng}},
2141 m_warnings{warnings},
2143 m_inbound_inv_bucket(m_opts.tx_send_rate, 1.0),
2148 if (opts.reconcile_txs) {
2153void PeerManagerImpl::StartScheduledTasks(
CScheduler& scheduler)
2164 scheduler.
scheduleFromNow([&] { ReattemptInitialBroadcast(scheduler); }, delta);
2166 if (m_opts.private_broadcast) {
2167 scheduler.
scheduleFromNow([&] { ReattemptPrivateBroadcast(scheduler); }, 0min);
2171void PeerManagerImpl::ActiveTipChange(
const CBlockIndex& new_tip,
bool is_ibd)
2179 LOCK(m_tx_download_mutex);
2183 m_txdownloadman.ActiveTipChange();
2193void PeerManagerImpl::BlockConnected(
2195 const std::shared_ptr<const CBlock>& pblock,
2200 m_last_tip_update = GetTime<std::chrono::seconds>();
2203 auto stalling_timeout = m_block_stalling_timeout.load();
2207 if (m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
2216 LOCK(m_tx_download_mutex);
2217 m_txdownloadman.BlockConnected(pblock);
2221void PeerManagerImpl::BlockDisconnected(
const std::shared_ptr<const CBlock> &block,
const CBlockIndex* pindex)
2223 LOCK(m_tx_download_mutex);
2224 m_txdownloadman.BlockDisconnected();
2231void PeerManagerImpl::NewPoWValidBlock(
const CBlockIndex *pindex,
const std::shared_ptr<const CBlock>& pblock)
2233 auto pcmpctblock = std::make_shared<const CBlockHeaderAndShortTxIDs>(*pblock,
FastRandomContext().rand64());
2237 if (pindex->
nHeight <= m_highest_fast_announce)
2239 m_highest_fast_announce = pindex->
nHeight;
2243 uint256 hashBlock(pblock->GetHash());
2244 const std::shared_future<CSerializedNetMsg> lazy_ser{
2248 auto most_recent_block_txs = std::make_unique<std::map<GenTxid, CTransactionRef>>();
2249 for (
const auto& tx : pblock->vtx) {
2250 most_recent_block_txs->emplace(tx->GetHash(), tx);
2251 most_recent_block_txs->emplace(tx->GetWitnessHash(), tx);
2254 LOCK(m_most_recent_block_mutex);
2255 m_most_recent_block_hash = hashBlock;
2256 m_most_recent_block = pblock;
2257 m_most_recent_compact_block = pcmpctblock;
2258 m_most_recent_block_txs = std::move(most_recent_block_txs);
2266 ProcessBlockAvailability(pnode->
GetId());
2270 if (state.m_requested_hb_cmpctblocks && !PeerHasHeader(&state, pindex) && PeerHasHeader(&state, pindex->
pprev)) {
2272 LogDebug(
BCLog::NET,
"%s sending header-and-ids %s to peer=%d\n",
"PeerManager::NewPoWValidBlock",
2273 hashBlock.ToString(), pnode->
GetId());
2276 PushMessage(*pnode, ser_cmpctblock.Copy());
2277 state.pindexBestHeaderSent = pindex;
2286void PeerManagerImpl::UpdatedBlockTip(
const CBlockIndex *pindexNew,
const CBlockIndex *pindexFork,
bool fInitialDownload)
2288 SetBestBlock(pindexNew->
nHeight, std::chrono::seconds{pindexNew->GetBlockTime()});
2291 if (fInitialDownload)
return;
2294 std::vector<uint256> vHashes;
2296 while (pindexToAnnounce != pindexFork) {
2298 pindexToAnnounce = pindexToAnnounce->
pprev;
2308 for (
auto& it : m_peer_map) {
2309 Peer& peer = *it.second;
2310 LOCK(peer.m_block_inv_mutex);
2311 for (
const uint256& hash : vHashes | std::views::reverse) {
2312 peer.m_blocks_for_headers_relay.push_back(hash);
2324void PeerManagerImpl::BlockChecked(
const std::shared_ptr<const CBlock>& block,
const BlockValidationState& state)
2328 const uint256 hash(block->GetHash());
2329 std::map<uint256, std::pair<NodeId, bool>>::iterator it = mapBlockSource.find(hash);
2334 it != mapBlockSource.end() &&
2335 State(it->second.first)) {
2336 MaybePunishNodeForBlock( it->second.first, state, !it->second.second);
2346 mapBlocksInFlight.count(hash) == mapBlocksInFlight.size()) {
2347 if (it != mapBlockSource.end()) {
2348 MaybeSetPeerAsAnnouncingHeaderAndIDs(it->second.first);
2351 if (it != mapBlockSource.end())
2352 mapBlockSource.erase(it);
2360bool PeerManagerImpl::AlreadyHaveBlock(
const uint256& block_hash)
2365void PeerManagerImpl::SendPings()
2368 for(
auto& it : m_peer_map) it.second->m_ping_queued =
true;
2371std::vector<Wtxid> InvToSendBucket::TakeForProcessing(
CTxMemPool& mempool)
2375 size_t n_to_take =
static_cast<size_t>(std::max<double>(count_bucket.
value() - count_floor, 0));
2377 std::vector<Wtxid> best;
2380 bool tokens_left =
true;
2381 for (
auto txiter : itervec) {
2382 auto& wtxid = txiter->GetTx().GetWitnessHash();
2384 best.push_back(wtxid);
2385 if (!decrement(txiter->GetTx().ComputeTotalSize())) {
2386 tokens_left =
false;
2389 backlog.push_back(wtxid);
2395 std::vector<Wtxid> dummy;
2397 dummy.swap(backlog);
2407 if (!backlog_bumped && now <= m_next_inv_bucket_check.load())
return;
2410 LOCK(m_inv_to_send_mutex);
2411 m_inbound_inv_bucket.increment(now);
2412 m_outbound_inv_bucket.increment(now);
2415 if (!m_next_inv_bucket_heartbeat.has_value()) {
2417 m_next_inv_bucket_heartbeat = now;
2420 if (m_next_inv_bucket_heartbeat.has_value() && now >= *m_next_inv_bucket_heartbeat) {
2421 LogDebug(
BCLog::NET,
"Transaction rate-limiting backlog inbound=%d itok=%.1f isz=%.1f outbound=%d otok=%.1f osz=%.1f",
2422 m_inbound_inv_bucket.backlog.size(),
2423 m_inbound_inv_bucket.count_bucket.value(),
2424 m_inbound_inv_bucket.size_bucket.value(),
2425 m_outbound_inv_bucket.backlog.size(),
2426 m_outbound_inv_bucket.count_bucket.value(),
2427 m_outbound_inv_bucket.size_bucket.value());
2428 if (m_inbound_inv_bucket.backlog.empty() && m_outbound_inv_bucket.backlog.empty()) {
2429 m_next_inv_bucket_heartbeat = std::nullopt;
2436 bool in_avail = m_inbound_inv_bucket.avail();
2437 bool out_avail = m_outbound_inv_bucket.avail();
2438 if (!in_avail && !out_avail)
return;
2440 std::vector<Wtxid> for_inbound;
2441 std::vector<Wtxid> for_outbound;
2445 if (in_avail) for_inbound = m_inbound_inv_bucket.TakeForProcessing(m_mempool);
2446 if (out_avail) for_outbound = m_outbound_inv_bucket.TakeForProcessing(m_mempool);
2449 if (!for_inbound.empty() || !for_outbound.empty()) {
2450 bool any_inbound_connected =
false;
2451 bool any_outbound_connected =
false;
2452 for (
const PeerRef& peer_ref : GetAllPeers()) {
2453 if (!peer_ref)
continue;
2454 Peer& peer{*peer_ref};
2455 auto tx_relay = peer.GetTxRelay();
2456 if (!tx_relay)
continue;
2458 LOCK(tx_relay->m_tx_inventory_mutex);
2464 if (tx_relay->m_next_inv_send_time == 0
s)
continue;
2465 if (peer.m_is_inbound) {
2466 any_inbound_connected =
true;
2468 any_outbound_connected =
true;
2470 for (
auto& i : (peer.m_is_inbound ? for_inbound : for_outbound)) {
2471 tx_relay->m_tx_inventory_to_send.push_back(i);
2478 if (!any_inbound_connected) m_inbound_inv_bucket.backlog.clear();
2479 if (!any_outbound_connected) m_outbound_inv_bucket.backlog.clear();
2483void PeerManagerImpl::InitiateTxBroadcastToAll(
const Wtxid& wtxid)
2486 LOCK(m_inv_to_send_mutex);
2487 m_inbound_inv_bucket.backlog.push_back(wtxid);
2488 m_outbound_inv_bucket.backlog.push_back(wtxid);
2495 const auto txstr{
strprintf(
"txid=%s, wtxid=%s", tx->GetHash().ToString(), tx->GetWitnessHash().ToString())};
2496 switch (m_tx_for_private_broadcast.Add(tx)) {
2511void PeerManagerImpl::RelayAddress(
NodeId originator,
2527 const auto current_time{GetTime<std::chrono::seconds>()};
2535 unsigned int nRelayNodes = (fReachable || (hasher.Finalize() & 1)) ? 2 : 1;
2537 std::array<std::pair<uint64_t, Peer*>, 2> best{{{0,
nullptr}, {0,
nullptr}}};
2538 assert(nRelayNodes <= best.size());
2542 for (
auto& [
id, peer] : m_peer_map) {
2543 if (peer->m_addr_relay_enabled &&
id != originator && IsAddrCompatible(*peer, addr)) {
2545 for (
unsigned int i = 0; i < nRelayNodes; i++) {
2546 if (hashKey > best[i].first) {
2547 std::copy(best.begin() + i, best.begin() + nRelayNodes - 1, best.begin() + i + 1);
2548 best[i] = std::make_pair(hashKey, peer.get());
2555 for (
unsigned int i = 0; i < nRelayNodes && best[i].first != 0; i++) {
2556 PushAddress(*best[i].second, addr);
2560void PeerManagerImpl::ProcessGetBlockData(
CNode& pfrom, Peer& peer,
const CInv& inv)
2570 std::shared_ptr<const CBlock> a_recent_block;
2571 std::shared_ptr<const CBlockHeaderAndShortTxIDs> a_recent_compact_block;
2573 LOCK(m_most_recent_block_mutex);
2574 a_recent_block = m_most_recent_block;
2575 a_recent_compact_block = m_most_recent_compact_block;
2578 bool need_activate_chain =
false;
2590 need_activate_chain =
true;
2594 if (need_activate_chain) {
2596 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
2603 bool can_direct_fetch{
false};
2611 if (!BlockRequestAllowed(*pindex)) {
2612 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block that isn't in the main chain\n", __func__, pfrom.
GetId());
2639 can_direct_fetch = CanDirectFetch();
2643 std::shared_ptr<const CBlock> pblock;
2644 if (a_recent_block && a_recent_block->GetHash() == inv.
hash) {
2645 pblock = a_recent_block;
2663 std::shared_ptr<CBlock> pblockRead = std::make_shared<CBlock>();
2673 pblock = pblockRead;
2681 bool sendMerkleBlock =
false;
2683 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
2684 LOCK(tx_relay->m_bloom_filter_mutex);
2685 if (tx_relay->m_bloom_filter) {
2686 sendMerkleBlock =
true;
2687 merkleBlock =
CMerkleBlock(*pblock, *tx_relay->m_bloom_filter);
2690 if (sendMerkleBlock) {
2698 for (
const auto& [tx_idx,
_] : merkleBlock.
vMatchedTxn)
2709 if (a_recent_compact_block && a_recent_compact_block->header.GetHash() == inv.
hash) {
2722 LOCK(peer.m_block_inv_mutex);
2724 if (inv.
hash == peer.m_continuation_block) {
2728 std::vector<CInv> vInv;
2729 vInv.emplace_back(
MSG_BLOCK, tip->GetBlockHash());
2731 peer.m_continuation_block.SetNull();
2739 auto txinfo{std::visit(
2740 [&](
const auto&
id) {
2741 return m_mempool.
info_for_relay(
id,
WITH_LOCK(tx_relay.m_tx_inventory_mutex,
return tx_relay.m_last_inv_sequence));
2745 return std::move(txinfo.tx);
2750 LOCK(m_most_recent_block_mutex);
2751 if (m_most_recent_block_txs !=
nullptr) {
2752 auto it = m_most_recent_block_txs->find(gtxid);
2753 if (it != m_most_recent_block_txs->end())
return it->second;
2760void PeerManagerImpl::ProcessGetData(
CNode& pfrom, Peer& peer,
const std::atomic<bool>& interruptMsgProc)
2764 auto tx_relay = peer.GetTxRelay();
2766 std::deque<CInv>::iterator it = peer.m_getdata_requests.begin();
2767 std::vector<CInv> vNotFound;
2772 while (it != peer.m_getdata_requests.end() && it->IsGenTxMsg()) {
2773 if (interruptMsgProc)
return;
2778 const CInv &inv = *it++;
2780 if (tx_relay ==
nullptr) {
2786 if (
auto tx{FindTxForGetData(*tx_relay,
ToGenTxid(inv))}) {
2789 MakeAndPushMessage(pfrom,
NetMsgType::TX, maybe_with_witness(*tx));
2792 vNotFound.push_back(inv);
2798 if (it != peer.m_getdata_requests.end() && !pfrom.
fPauseSend) {
2799 const CInv &inv = *it++;
2801 ProcessGetBlockData(pfrom, peer, inv);
2810 peer.m_getdata_requests.erase(peer.m_getdata_requests.begin(), it);
2812 if (!vNotFound.empty()) {
2831uint32_t PeerManagerImpl::GetFetchFlags(
const Peer& peer)
const
2833 uint32_t nFetchFlags = 0;
2834 if (CanServeWitnesses(peer)) {
2843 for (
size_t i = 0; i < req.
indexes.size(); i++) {
2845 Misbehaving(peer,
"getblocktxn with out-of-bounds tx indices");
2852 uint32_t tx_requested_size{0};
2853 for (
const auto& tx : resp.txn) tx_requested_size += tx->ComputeTotalSize();
2859bool PeerManagerImpl::CheckHeadersPoW(
const std::vector<CBlockHeader>&
headers, Peer& peer)
2863 Misbehaving(peer,
"header with invalid proof of work");
2868 if (!CheckHeadersAreContinuous(
headers)) {
2869 Misbehaving(peer,
"non-continuous headers sequence");
2894void PeerManagerImpl::HandleUnconnectingHeaders(
CNode& pfrom, Peer& peer,
2895 const std::vector<CBlockHeader>&
headers)
2899 if (MaybeSendGetHeaders(pfrom,
GetLocator(best_header), peer)) {
2900 LogDebug(
BCLog::NET,
"received header %s: missing prev block %s, sending getheaders (%d) to end (peer=%d)\n",
2902 headers[0].hashPrevBlock.ToString(),
2903 best_header->nHeight,
2913bool PeerManagerImpl::CheckHeadersAreContinuous(
const std::vector<CBlockHeader>&
headers)
const
2917 if (!hashLastBlock.
IsNull() && header.hashPrevBlock != hashLastBlock) {
2920 hashLastBlock = header.GetHash();
2925bool PeerManagerImpl::IsContinuationOfLowWorkHeadersSync(Peer& peer,
CNode& pfrom, std::vector<CBlockHeader>&
headers)
2927 if (peer.m_headers_sync) {
2928 auto result = peer.m_headers_sync->ProcessNextHeaders(
headers,
headers.size() == m_opts.max_headers_result);
2930 if (result.success) peer.m_last_getheaders_timestamp = {};
2931 if (result.request_more) {
2932 auto locator = peer.m_headers_sync->NextHeadersRequestLocator();
2934 Assume(!locator.vHave.empty());
2937 if (!locator.vHave.empty()) {
2940 bool sent_getheaders = MaybeSendGetHeaders(pfrom, locator, peer);
2943 locator.vHave.front().ToString(), pfrom.
GetId());
2948 peer.m_headers_sync.reset(
nullptr);
2953 LOCK(m_headers_presync_mutex);
2954 m_headers_presync_stats.erase(pfrom.
GetId());
2957 HeadersPresyncStats stats;
2958 stats.first = peer.m_headers_sync->GetPresyncWork();
2960 stats.second = {peer.m_headers_sync->GetPresyncHeight(),
2961 peer.m_headers_sync->GetPresyncTime()};
2965 LOCK(m_headers_presync_mutex);
2966 m_headers_presync_stats[pfrom.
GetId()] = stats;
2967 auto best_it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
2968 bool best_updated =
false;
2969 if (best_it == m_headers_presync_stats.end()) {
2973 const HeadersPresyncStats* stat_best{
nullptr};
2974 for (
const auto& [peer, stat] : m_headers_presync_stats) {
2975 if (!stat_best || stat > *stat_best) {
2980 m_headers_presync_bestpeer = peer_best;
2981 best_updated = (peer_best == pfrom.
GetId());
2982 }
else if (best_it->first == pfrom.
GetId() || stats > best_it->second) {
2984 m_headers_presync_bestpeer = pfrom.
GetId();
2985 best_updated =
true;
2987 if (best_updated && stats.second.has_value()) {
2989 m_headers_presync_should_signal =
true;
2993 if (result.success) {
2996 headers.swap(result.pow_validated_headers);
2999 return result.success;
3007bool PeerManagerImpl::TryLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
const CBlockIndex& chain_start_header, std::vector<CBlockHeader>&
headers)
3014 arith_uint256 minimum_chain_work = GetAntiDoSWorkThreshold();
3018 if (total_work < minimum_chain_work) {
3022 if (
headers.size() == m_opts.max_headers_result) {
3032 LOCK(peer.m_headers_sync_mutex);
3035 m_chainparams.
HeadersSync(), chain_start_header, minimum_chain_work));
3041 const auto msg{
strprintf(
"Failure when attempting to initiate headers sync: %s", e.what())};
3042 std::cerr <<
msg << std::endl;
3050 (void)IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3064bool PeerManagerImpl::IsAncestorOfBestHeaderOrTip(
const CBlockIndex* header)
3066 if (header ==
nullptr) {
3068 }
else if (m_chainman.m_best_header !=
nullptr && header == m_chainman.m_best_header->GetAncestor(header->
nHeight)) {
3076bool PeerManagerImpl::MaybeSendGetHeaders(
CNode& pfrom,
const CBlockLocator& locator, Peer& peer)
3084 peer.m_last_getheaders_timestamp = current_time;
3095void PeerManagerImpl::HeadersDirectFetchBlocks(
CNode& pfrom,
const Peer& peer,
const CBlockIndex& last_header)
3098 CNodeState *nodestate =
State(pfrom.
GetId());
3101 std::vector<const CBlockIndex*> vToFetch;
3109 vToFetch.push_back(pindexWalk);
3111 pindexWalk = pindexWalk->
pprev;
3123 std::vector<CInv> vGetData;
3125 for (
const CBlockIndex* pindex : vToFetch | std::views::reverse) {
3130 uint32_t nFetchFlags = GetFetchFlags(peer);
3132 BlockRequested(pfrom.
GetId(), *pindex);
3136 if (vGetData.size() > 1) {
3141 if (vGetData.size() > 0) {
3142 if (!m_opts.ignore_incoming_txs &&
3143 nodestate->m_provides_cmpctblocks &&
3144 vGetData.size() == 1 &&
3145 mapBlocksInFlight.size() == 1 &&
3161void PeerManagerImpl::UpdatePeerStateForReceivedHeaders(
CNode& pfrom,
3162 const CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
3165 CNodeState *nodestate =
State(pfrom.
GetId());
3174 nodestate->m_last_block_announcement =
GetTime();
3182 if (nodestate->pindexBestKnownBlock && nodestate->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
3204 if (m_outbound_peers_with_protect_from_disconnect < MAX_OUTBOUND_PEERS_TO_PROTECT_FROM_DISCONNECT && nodestate->pindexBestKnownBlock->nChainWork >= m_chainman.
ActiveChain().
Tip()->
nChainWork && !nodestate->m_chain_sync.m_protect) {
3206 nodestate->m_chain_sync.m_protect =
true;
3207 ++m_outbound_peers_with_protect_from_disconnect;
3212void PeerManagerImpl::ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
3213 std::vector<CBlockHeader>&&
headers,
3214 bool via_compact_block)
3216 size_t nCount =
headers.size();
3223 LOCK(peer.m_headers_sync_mutex);
3224 if (peer.m_headers_sync) {
3225 peer.m_headers_sync.reset(
nullptr);
3226 LOCK(m_headers_presync_mutex);
3227 m_headers_presync_stats.erase(pfrom.
GetId());
3231 peer.m_last_getheaders_timestamp = {};
3239 if (!CheckHeadersPoW(
headers, peer)) {
3254 bool already_validated_work =
false;
3257 bool have_headers_sync =
false;
3259 LOCK(peer.m_headers_sync_mutex);
3261 already_validated_work = IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3277 have_headers_sync = !!peer.m_headers_sync;
3282 bool headers_connect_blockindex{chain_start_header !=
nullptr};
3284 if (!headers_connect_blockindex) {
3288 HandleUnconnectingHeaders(pfrom, peer,
headers);
3295 peer.m_last_getheaders_timestamp = {};
3305 already_validated_work = already_validated_work || IsAncestorOfBestHeaderOrTip(last_received_header);
3312 already_validated_work =
true;
3318 if (!already_validated_work && TryLowWorkHeadersSync(peer, pfrom,
3319 *chain_start_header,
headers)) {
3331 bool received_new_header{last_received_header ==
nullptr};
3337 state, &pindexLast)};
3343 "If this happens with all peers, consider database corruption (that -reindex may fix) "
3344 "or a potential consensus incompatibility.",
3347 MaybePunishNodeForBlock(pfrom.
GetId(), state, via_compact_block,
"invalid header received");
3353 if (processed && received_new_header) {
3354 LogBlockHeader(*pindexLast, pfrom,
false);
3358 if (nCount == m_opts.max_headers_result && !have_headers_sync) {
3360 if (MaybeSendGetHeaders(pfrom,
GetLocator(pindexLast), peer)) {
3365 UpdatePeerStateForReceivedHeaders(pfrom, *pindexLast, received_new_header, nCount == m_opts.max_headers_result);
3368 HeadersDirectFetchBlocks(pfrom, peer, *pindexLast);
3374 bool first_time_failure)
3380 PeerRef peer{GetPeerRef(nodeid)};
3383 ptx->GetHash().ToString(),
3384 ptx->GetWitnessHash().ToString(),
3388 const auto& [add_extra_compact_tx, unique_parents, package_to_validate] = m_txdownloadman.MempoolRejectedTx(ptx, state, nodeid, first_time_failure);
3391 AddToCompactExtraTransactions(ptx);
3393 for (
const Txid& parent_txid : unique_parents) {
3394 if (peer) AddKnownTx(*peer, parent_txid.ToUint256());
3397 return package_to_validate;
3400void PeerManagerImpl::ProcessValidTx(
NodeId nodeid,
const CTransactionRef& tx,
const std::list<CTransactionRef>& replaced_transactions)
3406 m_txdownloadman.MempoolAcceptedTx(tx);
3410 tx->GetHash().ToString(),
3411 tx->GetWitnessHash().ToString(),
3414 InitiateTxBroadcastToAll(tx->GetWitnessHash());
3417 AddToCompactExtraTransactions(removedTx);
3427 const auto&
package = package_to_validate.m_txns;
3428 const auto& senders = package_to_validate.
m_senders;
3431 m_txdownloadman.MempoolRejectedPackage(package);
3434 if (!
Assume(package.size() == 2))
return;
3438 auto package_iter = package.rbegin();
3439 auto senders_iter = senders.rbegin();
3440 while (package_iter != package.rend()) {
3441 const auto& tx = *package_iter;
3442 const NodeId nodeid = *senders_iter;
3443 const auto it_result{package_result.
m_tx_results.find(tx->GetWitnessHash())};
3447 const auto& tx_result = it_result->second;
3448 switch (tx_result.m_result_type) {
3451 ProcessValidTx(nodeid, tx, tx_result.m_replaced_transactions);
3461 ProcessInvalidTx(nodeid, tx, tx_result.m_state,
false);
3479bool PeerManagerImpl::ProcessOrphanTx(Peer& peer)
3484 while (
CTransactionRef porphanTx = m_txdownloadman.GetTxToReconsider(peer.m_id)) {
3487 const Txid& orphanHash = porphanTx->GetHash();
3488 const Wtxid& orphan_wtxid = porphanTx->GetWitnessHash();
3505 ProcessInvalidTx(peer.m_id, porphanTx, state,
false);
3514bool PeerManagerImpl::PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
3516 const uint256& stop_hash, uint32_t max_height_diff,
3520 const bool supported_filter_type =
3523 if (!supported_filter_type) {
3525 static_cast<uint8_t
>(filter_type),
node.DisconnectMsg());
3526 node.fDisconnect =
true;
3535 if (!stop_index || !BlockRequestAllowed(*stop_index)) {
3538 node.fDisconnect =
true;
3543 uint32_t stop_height = stop_index->
nHeight;
3544 if (start_height > stop_height) {
3546 "start height %d and stop height %d, %s",
3547 start_height, stop_height,
node.DisconnectMsg());
3548 node.fDisconnect =
true;
3551 if (stop_height - start_height >= max_height_diff) {
3553 stop_height - start_height + 1, max_height_diff,
node.DisconnectMsg());
3554 node.fDisconnect =
true;
3559 if (!filter_index) {
3569 uint8_t filter_type_ser;
3570 uint32_t start_height;
3573 vRecv >> filter_type_ser >> start_height >> stop_hash;
3579 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3584 std::vector<BlockFilter> filters;
3586 LogDebug(
BCLog::NET,
"Failed to find block filter in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3591 for (
const auto& filter : filters) {
3598 uint8_t filter_type_ser;
3599 uint32_t start_height;
3602 vRecv >> filter_type_ser >> start_height >> stop_hash;
3608 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3614 if (start_height > 0) {
3616 stop_index->
GetAncestor(
static_cast<int>(start_height - 1));
3618 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3624 std::vector<uint256> filter_hashes;
3626 LogDebug(
BCLog::NET,
"Failed to find block filter hashes in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3640 uint8_t filter_type_ser;
3643 vRecv >> filter_type_ser >> stop_hash;
3649 if (!PrepareBlockFilterRequest(
node, peer, filter_type, 0, stop_hash,
3650 std::numeric_limits<uint32_t>::max(),
3651 stop_index, filter_index)) {
3659 for (
int i =
headers.size() - 1; i >= 0; i--) {
3664 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3678 bool new_block{
false};
3679 m_chainman.
ProcessNewBlock(block, force_processing, min_pow_checked, &new_block);
3681 node.m_last_block_time = GetTime<std::chrono::seconds>();
3686 RemoveBlockRequest(block->GetHash(), std::nullopt);
3689 mapBlockSource.erase(block->GetHash());
3693void PeerManagerImpl::ProcessCompactBlockTxns(
CNode& pfrom, Peer& peer,
const BlockTransactions& block_transactions)
3695 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
3696 bool fBlockRead{
false};
3700 auto range_flight = mapBlocksInFlight.equal_range(block_transactions.
blockhash);
3701 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
3702 bool requested_block_from_this_peer{
false};
3705 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
3707 while (range_flight.first != range_flight.second) {
3708 auto [node_id, block_it] = range_flight.first->second;
3709 if (node_id == pfrom.
GetId() && block_it->partialBlock) {
3710 requested_block_from_this_peer =
true;
3713 range_flight.first++;
3716 if (!requested_block_from_this_peer) {
3728 Misbehaving(peer,
"previous compact block reconstruction attempt failed");
3739 Misbehaving(peer,
"invalid compact block/non-matching block transactions");
3742 if (first_in_flight) {
3747 std::vector<CInv> invs;
3752 LogDebug(
BCLog::NET,
"Peer %d sent us a compact block but it failed to reconstruct, waiting on first download to complete\n", pfrom.
GetId());
3765 mapBlockSource.emplace(block_transactions.
blockhash, std::make_pair(pfrom.
GetId(),
false));
3780void PeerManagerImpl::LogBlockHeader(
const CBlockIndex& index,
const CNode& peer,
bool via_compact_block) {
3792 "Saw new %sheader hash=%s height=%d %s",
3793 via_compact_block ?
"cmpctblock " :
"",
3805void PeerManagerImpl::PushPrivateBroadcastTx(
CNode&
node)
3809 const auto opt_tx{m_tx_for_private_broadcast.PickTxForSend(
node.GetId(),
CService{node.addr})};
3812 node.fDisconnect =
true;
3818 tx->GetHash().ToString(), tx->HasWitness() ?
strprintf(
", wtxid=%s", tx->GetWitnessHash().ToString()) :
"",
3824void PeerManagerImpl::ProcessMessage(Peer& peer,
CNode& pfrom,
const std::string& msg_type,
DataStream& vRecv,
3826 const std::atomic<bool>& interruptMsgProc)
3841 uint64_t nNonce = 1;
3844 std::string cleanSubVer;
3845 int starting_height = -1;
3848 vRecv >> nVersion >> Using<CustomUintFormatter<8>>(nServices) >> nTime;
3863 LogDebug(
BCLog::NET,
"peer does not offer the expected services (%08x offered, %08x expected), %s",
3865 GetDesirableServiceFlags(nServices),
3878 if (!vRecv.
empty()) {
3886 if (!vRecv.
empty()) {
3887 std::string strSubVer;
3891 if (!vRecv.
empty()) {
3892 vRecv >> starting_height;
3912 PushNodeVersion(pfrom, peer);
3916 const int greatest_common_version = std::min(nVersion, pfrom.
AdvertisedVersion());
3921 peer.m_their_services = nServices;
3925 pfrom.cleanSubVer = cleanSubVer;
3936 (fRelay || (peer.m_our_services &
NODE_BLOOM))) {
3937 auto*
const tx_relay = peer.SetTxRelay();
3939 LOCK(tx_relay->m_bloom_filter_mutex);
3940 tx_relay->m_relay_txs = fRelay;
3946 LogDebug(
BCLog::NET,
"receive version message: %s: version %d, blocks=%d, us=%s, txrelay=%d, %s%s",
3947 cleanSubVer.empty() ?
"<no user agent>" : cleanSubVer, pfrom.
nVersion,
3949 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
3967 if (greatest_common_version >= 70016) {
3982 const auto* tx_relay = peer.GetTxRelay();
3983 if (tx_relay &&
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs) &&
3985 const uint64_t recon_salt = m_txreconciliation->PreRegisterPeer(pfrom.
GetId());
3998 if (MaybeDisconnectForTxRelayCapacity(pfrom, msg_type, pfrom.
GetId()))
return;
4006 m_num_preferred_download_peers += state->fPreferredDownload;
4012 bool send_getaddr{
false};
4014 send_getaddr = SetupAddressRelay(pfrom, peer);
4024 peer.m_getaddr_sent =
true;
4048 peer.m_time_offset =
NodeSeconds{std::chrono::seconds{nTime}} - Now<NodeSeconds>();
4052 m_outbound_time_offsets.Add(peer.m_time_offset);
4053 m_outbound_time_offsets.WarnIfOutOfSync();
4057 if (greatest_common_version <= 70012) {
4058 constexpr auto finalAlert{
"60010000000000000000000000ffffff7f00000000ffffff7ffeffff7f01ffffff7f00000000ffffff7f00ffffff7f002f555247454e543a20416c657274206b657920636f6d70726f6d697365642c2075706772616465207265717569726564004630440220653febd6410f470f6bae11cad19c48413becb1ac2c17f908fd0fd53bdc3abd5202206d0e9c96fe88d4a0f01ed9dedae2b6f9e00da94cad0fecaae66ecf689bf71b50"_hex};
4059 MakeAndPushMessage(pfrom,
"alert", finalAlert);
4082 auto new_peer_msg = [&]() {
4084 return strprintf(
"New %s peer connected: transport: %s, version: %d, %s%s",
4088 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
4096 LogInfo(
"%s", new_peer_msg());
4099 if (
auto tx_relay = peer.GetTxRelay()) {
4108 tx_relay->m_tx_inventory_mutex,
4109 return tx_relay->m_tx_inventory_to_send.empty() &&
4110 tx_relay->m_next_inv_send_time == 0
s));
4120 PushPrivateBroadcastTx(pfrom);
4133 if (m_txreconciliation) {
4134 if (!peer.m_wtxid_relay || !m_txreconciliation->IsPeerRegistered(pfrom.
GetId())) {
4138 m_txreconciliation->ForgetPeer(pfrom.
GetId());
4144 const CNodeState* state =
State(pfrom.
GetId());
4146 .m_preferred = state->fPreferredDownload,
4147 .m_relay_permissions = pfrom.HasPermission(NetPermissionFlags::Relay),
4148 .m_wtxid_relay = peer.m_wtxid_relay,
4157 peer.m_prefers_headers =
true;
4162 uint8_t sendcmpct_hb{0};
4163 uint64_t sendcmpct_version{0};
4164 vRecv >> sendcmpct_hb >> sendcmpct_version;
4168 if (sendcmpct_hb > 1) {
4169 Misbehaving(peer,
"invalid sendcmpct announce field");
4177 CNodeState* nodestate =
State(pfrom.
GetId());
4178 nodestate->m_provides_cmpctblocks =
true;
4179 nodestate->m_requested_hb_cmpctblocks = sendcmpct_hb;
4196 if (!peer.m_wtxid_relay) {
4197 peer.m_wtxid_relay =
true;
4198 m_wtxid_relay_peers++;
4217 peer.m_wants_addrv2 =
true;
4234 std::string feature_id;
4238 std::vector<unsigned char> feature_data_vec;
4241 }
catch (
const std::exception&) {
4244 if (feature_id.size() < 4 || !vRecv.
empty()) {
4264 if (!m_txreconciliation) {
4265 LogDebug(
BCLog::NET,
"sendtxrcncl from peer=%d ignored, as our node does not have txreconciliation enabled\n", pfrom.
GetId());
4276 if (RejectIncomingTxs(pfrom)) {
4285 const auto* tx_relay = peer.GetTxRelay();
4286 if (!tx_relay || !
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs)) {
4292 uint32_t peer_txreconcl_version;
4293 uint64_t remote_salt;
4294 vRecv >> peer_txreconcl_version >> remote_salt;
4297 peer_txreconcl_version, remote_salt);
4329 const auto ser_params{
4337 std::vector<CAddress> vAddr;
4338 vRecv >> ser_params(vAddr);
4339 ProcessAddrs(msg_type, pfrom, peer, std::move(vAddr), interruptMsgProc);
4344 std::vector<CInv> vInv;
4348 Misbehaving(peer,
strprintf(
"inv message size = %u", vInv.size()));
4352 const bool reject_tx_invs{RejectIncomingTxs(pfrom)};
4353 std::unordered_set<uint256, SaltedUint256Hasher> seen_txids{0, m_txhash_hasher};
4354 std::unordered_set<uint256, SaltedUint256Hasher> seen_wtxids{0, m_txhash_hasher};
4358 const auto current_time{GetTime<std::chrono::microseconds>()};
4361 for (
CInv& inv : vInv) {
4362 if (interruptMsgProc)
return;
4367 if (peer.m_wtxid_relay) {
4374 const bool fAlreadyHave = AlreadyHaveBlock(inv.
hash);
4377 UpdateBlockAvailability(pfrom.
GetId(), inv.
hash);
4385 best_block = &inv.
hash;
4388 if (reject_tx_invs) {
4394 auto& seen_hashes{inv.
IsMsgWtx() ? seen_wtxids : seen_txids};
4395 if (!seen_hashes.insert(inv.
hash).second)
continue;
4397 AddKnownTx(peer, inv.
hash);
4400 const bool fAlreadyHave{m_txdownloadman.AddTxAnnouncement(pfrom.
GetId(), gtxid, current_time)};
4408 if (best_block !=
nullptr) {
4420 if (state.fSyncStarted || (!peer.m_inv_triggered_getheaders_before_sync && *best_block != m_last_block_inv_triggering_headers_sync)) {
4421 if (MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer)) {
4423 m_chainman.m_best_header->nHeight, best_block->ToString(),
4426 if (!state.fSyncStarted) {
4427 peer.m_inv_triggered_getheaders_before_sync =
true;
4431 m_last_block_inv_triggering_headers_sync = *best_block;
4440 std::vector<CInv> vInv;
4444 Misbehaving(peer,
strprintf(
"getdata message size = %u", vInv.size()));
4450 if (vInv.size() > 0) {
4455 const auto pushed_tx_opt{m_tx_for_private_broadcast.GetTxForNode(pfrom.
GetId())};
4456 if (!pushed_tx_opt) {
4467 if (vInv.size() == 1 && vInv[0].IsMsgTx() && vInv[0].hash == pushed_tx->GetHash().ToUint256()) {
4471 peer.m_ping_queued =
true;
4482 LOCK(peer.m_getdata_requests_mutex);
4483 peer.m_getdata_requests.insert(peer.m_getdata_requests.end(), vInv.begin(), vInv.end());
4484 ProcessGetData(pfrom, peer, interruptMsgProc);
4493 vRecv >> locator >> hashStop;
4509 std::shared_ptr<const CBlock> a_recent_block;
4511 LOCK(m_most_recent_block_mutex);
4512 a_recent_block = m_most_recent_block;
4515 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
4530 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4545 if (--nLimit <= 0) {
4549 WITH_LOCK(peer.m_block_inv_mutex, {peer.m_continuation_block = pindex->GetBlockHash();});
4569 for (
size_t i = 1; i < req.
indexes.size(); ++i) {
4573 std::shared_ptr<const CBlock> recent_block;
4575 LOCK(m_most_recent_block_mutex);
4576 if (m_most_recent_block_hash == req.
blockhash)
4577 recent_block = m_most_recent_block;
4581 SendBlockTransactions(pfrom, peer, *recent_block, req);
4600 if (!block_pos.IsNull()) {
4607 SendBlockTransactions(pfrom, peer, block, req);
4620 WITH_LOCK(peer.m_getdata_requests_mutex, peer.m_getdata_requests.push_back(inv));
4628 vRecv >> locator >> hashStop;
4646 if (m_chainman.
ActiveTip() ==
nullptr ||
4648 LogDebug(
BCLog::NET,
"Ignoring getheaders from peer=%d because active chain has too little work; sending empty response\n", pfrom.
GetId());
4655 CNodeState *nodestate =
State(pfrom.
GetId());
4664 if (!BlockRequestAllowed(*pindex)) {
4665 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block header that isn't in the main chain\n", __func__, pfrom.
GetId());
4678 std::vector<CBlock> vHeaders;
4679 int nLimit = m_opts.max_headers_result;
4681 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4684 if (--nLimit <= 0 || pindex->GetBlockHash() == hashStop)
4699 nodestate->pindexBestHeaderSent = pindex ? pindex : m_chainman.
ActiveChain().
Tip();
4705 if (RejectIncomingTxs(pfrom)) {
4719 const Txid& txid = ptx->GetHash();
4720 const Wtxid& wtxid = ptx->GetWitnessHash();
4723 AddKnownTx(peer, hash);
4725 if (
const auto num_broadcasted{m_tx_for_private_broadcast.Remove(ptx)}) {
4727 "network from %s; stopping private broadcast attempts",
4738 const auto& [should_validate, package_to_validate] = m_txdownloadman.ReceivedTx(pfrom.
GetId(), ptx);
4739 if (!should_validate) {
4744 if (!m_mempool.
exists(txid)) {
4745 LogInfo(
"Not relaying non-mempool transaction %s (wtxid=%s) from forcerelay peer=%d\n",
4748 LogInfo(
"Force relaying tx %s (wtxid=%s) from peer=%d\n",
4750 InitiateTxBroadcastToAll(wtxid);
4754 if (package_to_validate) {
4757 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4758 ProcessPackageResult(package_to_validate.value(), package_result);
4764 Assume(!package_to_validate.has_value());
4774 if (
auto package_to_validate{ProcessInvalidTx(pfrom.
GetId(), ptx, state,
true)}) {
4777 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4778 ProcessPackageResult(package_to_validate.value(), package_result);
4791 }
else if (m_opts.ignore_incoming_txs) {
4798 const CNodeState *nodestate =
State(pfrom.
GetId());
4799 if (!nodestate->m_provides_cmpctblocks) {
4806 vRecv >> cmpctblock;
4808 bool received_new_header =
false;
4818 MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer);
4828 received_new_header =
true;
4836 MaybePunishNodeForBlock(pfrom.
GetId(), state,
true,
"invalid header via cmpctblock");
4843 if (received_new_header) {
4844 LogBlockHeader(*pindex, pfrom,
true);
4847 bool fProcessBLOCKTXN =
false;
4851 bool fRevertToHeaderProcessing =
false;
4855 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
4856 bool fBlockReconstructed =
false;
4862 CNodeState *nodestate =
State(pfrom.
GetId());
4867 nodestate->m_last_block_announcement =
GetTime();
4873 auto range_flight = mapBlocksInFlight.equal_range(pindex->
GetBlockHash());
4874 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
4875 bool requested_block_from_this_peer{
false};
4878 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
4880 while (range_flight.first != range_flight.second) {
4881 if (range_flight.first->second.first == pfrom.
GetId()) {
4882 requested_block_from_this_peer =
true;
4885 range_flight.first++;
4895 if (requested_block_from_this_peer) {
4898 std::vector<CInv> vInv(1);
4906 if (!already_in_flight && !CanDirectFetch()) {
4914 requested_block_from_this_peer) {
4915 std::list<QueuedBlock>::iterator* queuedBlockIt =
nullptr;
4916 if (!BlockRequested(pfrom.
GetId(), *pindex, &queuedBlockIt)) {
4917 if (!(*queuedBlockIt)->partialBlock)
4930 Misbehaving(peer,
"invalid compact block");
4933 if (first_in_flight) {
4935 std::vector<CInv> vInv(1);
4946 for (
size_t i = 0; i < cmpctblock.
BlockTxCount(); i++) {
4951 fProcessBLOCKTXN =
true;
4952 }
else if (first_in_flight) {
4959 IsBlockRequestedFromOutbound(blockhash) ||
4978 ReadStatus status = tempBlock.InitData(cmpctblock, vExtraTxnForCompact);
4983 std::vector<CTransactionRef> dummy;
4985 status = tempBlock.FillBlock(*pblock, dummy,
4988 fBlockReconstructed =
true;
4992 if (requested_block_from_this_peer) {
4995 std::vector<CInv> vInv(1);
5001 fRevertToHeaderProcessing =
true;
5006 if (fProcessBLOCKTXN) {
5009 return ProcessCompactBlockTxns(pfrom, peer, txn);
5012 if (fRevertToHeaderProcessing) {
5018 return ProcessHeadersMessage(pfrom, peer, {cmpctblock.
header},
true);
5021 if (fBlockReconstructed) {
5026 mapBlockSource.emplace(pblock->GetHash(), std::make_pair(pfrom.
GetId(),
false));
5044 RemoveBlockRequest(pblock->GetHash(), std::nullopt);
5061 return ProcessCompactBlockTxns(pfrom, peer, resp);
5072 std::vector<CBlockHeader>
headers;
5076 if (nCount > m_opts.max_headers_result) {
5077 Misbehaving(peer,
strprintf(
"headers message size = %u", nCount));
5081 for (
unsigned int n = 0; n < nCount; n++) {
5086 ProcessHeadersMessage(pfrom, peer, std::move(
headers),
false);
5090 if (m_headers_presync_should_signal.exchange(
false)) {
5091 HeadersPresyncStats stats;
5093 LOCK(m_headers_presync_mutex);
5094 auto it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
5095 if (it != m_headers_presync_stats.end()) stats = it->second;
5113 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
5124 Misbehaving(peer,
"mutated block");
5129 bool forceProcessing =
false;
5130 const uint256 hash(pblock->GetHash());
5131 bool min_pow_checked =
false;
5136 forceProcessing = IsBlockRequested(hash);
5137 RemoveBlockRequest(hash, pfrom.
GetId());
5141 mapBlockSource.emplace(hash, std::make_pair(pfrom.
GetId(),
true));
5145 min_pow_checked =
true;
5148 ProcessBlock(pfrom, pblock, forceProcessing, min_pow_checked);
5165 Assume(SetupAddressRelay(pfrom, peer));
5169 if (peer.m_getaddr_recvd) {
5173 peer.m_getaddr_recvd =
true;
5175 peer.m_addrs_to_send.clear();
5176 std::vector<CAddress> vAddr;
5182 for (
const CAddress &addr : vAddr) {
5183 PushAddress(peer, addr);
5211 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5212 LOCK(tx_relay->m_tx_inventory_mutex);
5213 tx_relay->m_send_mempool =
true;
5239 ProcessPong(pfrom, peer, time_received, vRecv);
5255 Misbehaving(peer,
"too-large bloom filter");
5256 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5258 LOCK(tx_relay->m_bloom_filter_mutex);
5259 tx_relay->m_bloom_filter.reset(
new CBloomFilter(filter));
5260 tx_relay->m_relay_txs =
true;
5264 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5275 std::vector<unsigned char> vData;
5283 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5284 LOCK(tx_relay->m_bloom_filter_mutex);
5285 if (tx_relay->m_bloom_filter) {
5286 tx_relay->m_bloom_filter->insert(vData);
5292 Misbehaving(peer,
"bad filteradd message");
5303 auto tx_relay = peer.GetTxRelay();
5304 if (!tx_relay)
return;
5307 LOCK(tx_relay->m_bloom_filter_mutex);
5308 tx_relay->m_bloom_filter =
nullptr;
5309 tx_relay->m_relay_txs =
true;
5313 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5319 vRecv >> newFeeFilter;
5321 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5322 tx_relay->m_fee_filter_received = newFeeFilter;
5330 ProcessGetCFilters(pfrom, peer, vRecv);
5335 ProcessGetCFHeaders(pfrom, peer, vRecv);
5340 ProcessGetCFCheckPt(pfrom, peer, vRecv);
5345 std::vector<CInv> vInv;
5347 std::vector<GenTxid> tx_invs;
5349 for (
CInv &inv : vInv) {
5355 LOCK(m_tx_download_mutex);
5356 m_txdownloadman.ReceivedNotFound(pfrom.
GetId(), tx_invs);
5365bool PeerManagerImpl::MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer)
5368 LOCK(peer.m_misbehavior_mutex);
5371 if (!peer.m_should_discourage)
return false;
5373 peer.m_should_discourage =
false;
5378 LogWarning(
"Not punishing noban peer %d!", peer.m_id);
5384 LogWarning(
"Not punishing manually connected peer %d!", peer.m_id);
5404bool PeerManagerImpl::MaybeDisconnectForTxRelayCapacity(
CNode&
node,
const std::string& msg_type, std::optional<NodeId> protect_peer)
5406 if (!
node.IsInboundConn() || !
node.m_relays_txs)
return false;
5409 LogDebug(
BCLog::NET,
"failed to find a tx-relaying eviction candidate - connection dropped after %s message, peer=%d\n", msg_type,
node.GetId());
5410 node.fDisconnect =
true;
5414bool PeerManagerImpl::ProcessMessages(
CNode&
node, std::atomic<bool>& interruptMsgProc)
5419 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
5420 if (maybe_peer ==
nullptr)
return false;
5421 Peer& peer{*maybe_peer};
5425 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent)
return false;
5428 LOCK(peer.m_getdata_requests_mutex);
5429 if (!peer.m_getdata_requests.empty()) {
5430 ProcessGetData(
node, peer, interruptMsgProc);
5434 const bool processed_orphan = ProcessOrphanTx(peer);
5436 if (
node.fDisconnect)
5439 if (processed_orphan)
return true;
5444 LOCK(peer.m_getdata_requests_mutex);
5445 if (!peer.m_getdata_requests.empty())
return true;
5449 if (
node.fPauseSend)
return false;
5451 auto poll_result{
node.PollMessage()};
5458 bool fMoreWork = poll_result->second;
5462 node.m_addr_name.c_str(),
5463 node.ConnectionTypeAsString().c_str(),
5469 if (m_opts.capture_messages) {
5474 ProcessMessage(peer,
node,
msg.m_type,
msg.m_recv,
msg.m_time, interruptMsgProc);
5475 if (interruptMsgProc)
return false;
5477 LOCK(peer.m_getdata_requests_mutex);
5478 if (!peer.m_getdata_requests.empty()) fMoreWork =
true;
5485 LOCK(m_tx_download_mutex);
5486 if (m_txdownloadman.HaveMoreWork(peer.m_id)) fMoreWork =
true;
5487 }
catch (
const std::exception& e) {
5496void PeerManagerImpl::ConsiderEviction(
CNode& pto, Peer& peer, std::chrono::seconds time_in_seconds)
5509 if (state.pindexBestKnownBlock !=
nullptr && state.pindexBestKnownBlock->nChainWork >= m_chainman.
ActiveChain().
Tip()->
nChainWork) {
5511 if (state.m_chain_sync.m_timeout != 0
s) {
5512 state.m_chain_sync.m_timeout = 0
s;
5513 state.m_chain_sync.m_work_header =
nullptr;
5514 state.m_chain_sync.m_sent_getheaders =
false;
5516 }
else if (state.m_chain_sync.m_timeout == 0
s || (state.m_chain_sync.m_work_header !=
nullptr && state.pindexBestKnownBlock !=
nullptr && state.pindexBestKnownBlock->nChainWork >= state.m_chain_sync.m_work_header->nChainWork)) {
5524 state.m_chain_sync.m_work_header = m_chainman.
ActiveChain().
Tip();
5525 state.m_chain_sync.m_sent_getheaders =
false;
5526 }
else if (state.m_chain_sync.m_timeout > 0
s && time_in_seconds > state.m_chain_sync.m_timeout) {
5530 if (state.m_chain_sync.m_sent_getheaders) {
5532 LogInfo(
"Outbound peer has old chain, best known block = %s, %s", state.pindexBestKnownBlock !=
nullptr ? state.pindexBestKnownBlock->GetBlockHash().ToString() :
"<none>", pto.
DisconnectMsg());
5535 assert(state.m_chain_sync.m_work_header);
5540 MaybeSendGetHeaders(pto,
5541 GetLocator(state.m_chain_sync.m_work_header->pprev),
5543 LogDebug(
BCLog::NET,
"sending getheaders to outbound peer=%d to verify chain work (current best known block:%s, benchmark blockhash: %s)\n", pto.
GetId(), state.pindexBestKnownBlock !=
nullptr ? state.pindexBestKnownBlock->GetBlockHash().ToString() :
"<none>", state.m_chain_sync.m_work_header->GetBlockHash().ToString());
5544 state.m_chain_sync.m_sent_getheaders =
true;
5565 std::pair<NodeId, std::chrono::seconds> youngest_peer{-1, 0}, next_youngest_peer{-1, 0};
5569 if (pnode->
GetId() > youngest_peer.first) {
5570 next_youngest_peer = youngest_peer;
5571 youngest_peer.first = pnode->GetId();
5572 youngest_peer.second = pnode->m_last_block_time;
5575 NodeId to_disconnect = youngest_peer.first;
5576 if (youngest_peer.second > next_youngest_peer.second) {
5579 to_disconnect = next_youngest_peer.first;
5588 CNodeState *node_state =
State(pnode->
GetId());
5589 if (node_state ==
nullptr ||
5592 LogDebug(
BCLog::NET,
"disconnecting extra block-relay-only peer=%d (last block received at time %d)\n",
5596 LogDebug(
BCLog::NET,
"keeping block-relay-only peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5597 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), node_state->vBlocksInFlight.size());
5612 int64_t oldest_block_announcement = std::numeric_limits<int64_t>::max();
5615 AssertLockHeld(::cs_main);
5619 if (!pnode->IsFullOutboundConn() || pnode->fDisconnect) return;
5620 CNodeState *state = State(pnode->GetId());
5621 if (state == nullptr) return;
5623 if (state->m_chain_sync.m_protect) return;
5626 if (!m_connman.MultipleManualOrFullOutboundConns(pnode->addr.GetNetwork())) return;
5627 if (state->m_last_block_announcement < oldest_block_announcement || (state->m_last_block_announcement == oldest_block_announcement && pnode->GetId() > worst_peer)) {
5628 worst_peer = pnode->GetId();
5629 oldest_block_announcement = state->m_last_block_announcement;
5632 if (worst_peer != -1) {
5643 LogDebug(
BCLog::NET,
"disconnecting extra outbound peer=%d (last block announcement received at time %d)\n", pnode->
GetId(), oldest_block_announcement);
5647 LogDebug(
BCLog::NET,
"keeping outbound peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5648 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), state.vBlocksInFlight.size());
5664void PeerManagerImpl::CheckForStaleTipAndEvictPeers()
5669 auto now{GetTime<std::chrono::seconds>()};
5671 EvictExtraOutboundPeers(current_time);
5673 if (now > m_stale_tip_check_time) {
5677 LogInfo(
"Potential stale tip detected, will try using extra outbound peer (last tip update: %d seconds ago)\n",
5686 if (!m_initial_sync_finished && CanDirectFetch()) {
5688 m_initial_sync_finished =
true;
5695 peer.m_ping_nonce_sent &&
5705 bool pingSend =
false;
5707 if (peer.m_ping_queued) {
5712 if (peer.m_ping_nonce_sent == 0 && now > peer.m_ping_start.load() +
PING_INTERVAL) {
5721 }
while (
nonce == 0);
5722 peer.m_ping_queued =
false;
5723 peer.m_ping_start = now;
5725 peer.m_ping_nonce_sent =
nonce;
5729 peer.m_ping_nonce_sent = 0;
5735void PeerManagerImpl::MaybeSendAddr(
CNode&
node, Peer& peer, std::chrono::microseconds current_time)
5738 if (!peer.m_addr_relay_enabled)
return;
5740 LOCK(peer.m_addr_send_times_mutex);
5743 peer.m_next_local_addr_send < current_time) {
5750 if (peer.m_next_local_addr_send != 0us) {
5751 peer.m_addr_known->reset();
5754 CAddress local_addr{*local_service, peer.m_our_services, Now<NodeSeconds>()};
5755 if (peer.m_next_local_addr_send == 0us) {
5759 if (IsAddrCompatible(peer, local_addr)) {
5760 std::vector<CAddress> self_announcement{local_addr};
5761 if (peer.m_wants_addrv2) {
5769 PushAddress(peer, local_addr);
5776 if (current_time <= peer.m_next_addr_send)
return;
5789 bool ret = peer.m_addr_known->contains(addr.
GetKey());
5790 if (!
ret) peer.m_addr_known->insert(addr.
GetKey());
5793 peer.m_addrs_to_send.erase(std::remove_if(peer.m_addrs_to_send.begin(), peer.m_addrs_to_send.end(), addr_already_known),
5794 peer.m_addrs_to_send.end());
5797 if (peer.m_addrs_to_send.empty())
return;
5799 if (peer.m_wants_addrv2) {
5804 peer.m_addrs_to_send.clear();
5807 if (peer.m_addrs_to_send.capacity() > 40) {
5808 peer.m_addrs_to_send.shrink_to_fit();
5812void PeerManagerImpl::MaybeSendSendHeaders(
CNode&
node, Peer& peer)
5820 CNodeState &state = *
State(
node.GetId());
5821 if (state.pindexBestKnownBlock !=
nullptr &&
5828 peer.m_sent_sendheaders =
true;
5833void PeerManagerImpl::MaybeSendFeefilter(
CNode& pto, Peer& peer, std::chrono::microseconds current_time)
5835 if (m_opts.ignore_incoming_txs)
return;
5851 if (peer.m_fee_filter_sent == MAX_FILTER) {
5854 peer.m_next_send_feefilter = 0us;
5857 if (current_time > peer.m_next_send_feefilter) {
5858 CAmount filterToSend = m_fee_filter_rounder.round(currentFilter);
5861 if (filterToSend != peer.m_fee_filter_sent) {
5863 peer.m_fee_filter_sent = filterToSend;
5870 (currentFilter < 3 * peer.m_fee_filter_sent / 4 || currentFilter > 4 * peer.m_fee_filter_sent / 3)) {
5875bool PeerManagerImpl::RejectIncomingTxs(
const CNode& peer)
const
5888 const size_t nAvail{vRecv.
size()};
5889 bool bPingFinished =
false;
5890 std::string sProblem;
5892 if (nAvail >=
sizeof(
nonce)) {
5896 if (peer.m_ping_nonce_sent != 0) {
5897 if (
nonce == peer.m_ping_nonce_sent) {
5899 bPingFinished =
true;
5900 const auto ping_time = ping_end - peer.m_ping_start.load();
5901 if (ping_time.count() >= 0) {
5905 m_tx_for_private_broadcast.NodeConfirmedReception(pfrom.
GetId());
5912 sProblem =
"Timing mishap";
5916 sProblem =
"Nonce mismatch";
5919 bPingFinished =
true;
5920 sProblem =
"Nonce zero";
5924 sProblem =
"Unsolicited pong without ping";
5928 bPingFinished =
true;
5929 sProblem =
"Short payload";
5932 if (!(sProblem.empty())) {
5936 peer.m_ping_nonce_sent,
5940 if (bPingFinished) {
5941 peer.m_ping_nonce_sent = 0;
5945bool PeerManagerImpl::SetupAddressRelay(
const CNode&
node, Peer& peer)
5950 if (
node.IsBlockOnlyConn())
return false;
5955 if (
node.IsFeelerConn())
return false;
5957 if (!peer.m_addr_relay_enabled.exchange(
true)) {
5961 peer.m_addr_known = std::make_unique<CRollingBloomFilter>(5000, 0.001);
5967void PeerManagerImpl::ProcessAddrs(std::string_view msg_type,
CNode& pfrom, Peer& peer, std::vector<CAddress>&& vAddr,
const std::atomic<bool>& interruptMsgProc)
5972 if (!SetupAddressRelay(pfrom, peer)) {
5979 Misbehaving(peer,
strprintf(
"%s message size = %u", msg_type, vAddr.size()));
5984 std::vector<CAddress> vAddrOk;
5990 const auto time_diff{current_time - peer.m_addr_token_timestamp};
5994 peer.m_addr_token_timestamp = current_time;
5997 uint64_t num_proc = 0;
5998 uint64_t num_rate_limit = 0;
5999 std::shuffle(vAddr.begin(), vAddr.end(), m_rng);
6002 if (interruptMsgProc)
6006 if (peer.m_addr_token_bucket < 1.0) {
6012 peer.m_addr_token_bucket -= 1.0;
6021 addr.
nTime = std::chrono::time_point_cast<std::chrono::seconds>(current_time - 5 * 24h);
6023 AddAddressKnown(peer, addr);
6030 if (addr.
nTime > current_time - 10min && !peer.m_getaddr_sent && vAddr.size() <= 10 && addr.
IsRoutable()) {
6032 RelayAddress(pfrom.
GetId(), addr, reachable);
6036 vAddrOk.push_back(addr);
6039 peer.m_addr_processed += num_proc;
6040 peer.m_addr_rate_limited += num_rate_limit;
6041 LogDebug(
BCLog::NET,
"Received addr: %u addresses (%u processed, %u rate-limited) from peer=%d\n",
6042 vAddr.size(), num_proc, num_rate_limit, pfrom.
GetId());
6044 m_addrman.
Add(vAddrOk, pfrom.
addr, 2h);
6045 if (vAddr.size() < 1000) peer.m_getaddr_sent =
false;
6054bool PeerManagerImpl::SendMessages(
CNode&
node)
6059 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
6060 if (!maybe_peer)
return false;
6061 Peer& peer{*maybe_peer};
6066 if (MaybeDiscourageAndDisconnect(
node, peer))
return true;
6069 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent) {
6070 PushNodeVersion(
node, peer);
6071 peer.m_outbound_version_message_sent =
true;
6075 if (!
node.fSuccessfullyConnected ||
node.fDisconnect)
6079 const auto current_time{GetTime<std::chrono::microseconds>()};
6084 if (
node.IsPrivateBroadcastConn()) {
6088 node.fDisconnect =
true;
6095 node.fDisconnect =
true;
6099 MaybeSendPing(
node, peer, now);
6102 if (
node.fDisconnect)
return true;
6104 MaybeSendAddr(
node, peer, current_time);
6106 MaybeSendSendHeaders(
node, peer);
6108 ProcessInvBacklog(now);
6113 CNodeState &state = *
State(
node.GetId());
6116 if (m_chainman.m_best_header ==
nullptr) {
6123 bool sync_blocks_and_headers_from_peer =
false;
6124 if (state.fPreferredDownload) {
6125 sync_blocks_and_headers_from_peer =
true;
6126 }
else if (CanServeBlocks(peer) && !
node.IsAddrFetchConn()) {
6136 if (m_num_preferred_download_peers == 0 || mapBlocksInFlight.empty()) {
6137 sync_blocks_and_headers_from_peer =
true;
6143 if ((nSyncStarted == 0 && sync_blocks_and_headers_from_peer) || m_chainman.m_best_header->Time() >
NodeClock::now() - 24h) {
6144 const CBlockIndex* pindexStart = m_chainman.m_best_header;
6152 if (pindexStart->
pprev)
6153 pindexStart = pindexStart->
pprev;
6157 state.fSyncStarted =
true;
6181 LOCK(peer.m_block_inv_mutex);
6182 std::vector<CBlock> vHeaders;
6183 bool fRevertToInv = ((!peer.m_prefers_headers &&
6184 (!state.m_requested_hb_cmpctblocks || peer.m_blocks_for_headers_relay.size() > 1)) ||
6187 ProcessBlockAvailability(
node.GetId());
6189 if (!fRevertToInv) {
6190 bool fFoundStartingHeader =
false;
6194 for (
const uint256& hash : peer.m_blocks_for_headers_relay) {
6199 fRevertToInv =
true;
6202 if (pBestIndex !=
nullptr && pindex->
pprev != pBestIndex) {
6214 fRevertToInv =
true;
6217 pBestIndex = pindex;
6218 if (fFoundStartingHeader) {
6221 }
else if (PeerHasHeader(&state, pindex)) {
6223 }
else if (pindex->
pprev ==
nullptr || PeerHasHeader(&state, pindex->
pprev)) {
6226 fFoundStartingHeader =
true;
6231 fRevertToInv =
true;
6236 if (!fRevertToInv && !vHeaders.empty()) {
6237 if (vHeaders.size() == 1 && state.m_requested_hb_cmpctblocks) {
6241 vHeaders.front().GetHash().ToString(),
node.GetId());
6243 std::optional<CSerializedNetMsg> cached_cmpctblock_msg;
6245 LOCK(m_most_recent_block_mutex);
6246 if (m_most_recent_block_hash == pBestIndex->
GetBlockHash()) {
6250 if (cached_cmpctblock_msg.has_value()) {
6251 PushMessage(
node, std::move(cached_cmpctblock_msg.value()));
6259 state.pindexBestHeaderSent = pBestIndex;
6260 }
else if (peer.m_prefers_headers) {
6261 if (vHeaders.size() > 1) {
6264 vHeaders.front().GetHash().ToString(),
6265 vHeaders.back().GetHash().ToString(),
node.GetId());
6268 vHeaders.front().GetHash().ToString(),
node.GetId());
6271 state.pindexBestHeaderSent = pBestIndex;
6273 fRevertToInv =
true;
6279 if (!peer.m_blocks_for_headers_relay.empty()) {
6280 const uint256& hashToAnnounce = peer.m_blocks_for_headers_relay.back();
6293 if (!PeerHasHeader(&state, pindex)) {
6294 peer.m_blocks_for_inv_relay.push_back(hashToAnnounce);
6300 peer.m_blocks_for_headers_relay.clear();
6306 std::vector<CInv> vInv;
6308 LOCK(peer.m_block_inv_mutex);
6309 vInv.reserve(peer.m_blocks_for_inv_relay.size());
6312 for (
const uint256& hash : peer.m_blocks_for_inv_relay) {
6319 peer.m_blocks_for_inv_relay.clear();
6322 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
6323 LOCK(tx_relay->m_tx_inventory_mutex);
6326 if (tx_relay->m_next_inv_send_time < current_time) {
6327 fSendTrickle =
true;
6328 if (
node.IsInboundConn()) {
6337 LOCK(tx_relay->m_bloom_filter_mutex);
6338 if (!tx_relay->m_relay_txs) tx_relay->m_tx_inventory_to_send.clear();
6342 if (fSendTrickle && tx_relay->m_send_mempool) {
6343 auto vtxinfo = m_mempool.
infoAll();
6348 tx_relay->m_send_mempool =
false;
6349 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6352 tx_relay->m_tx_inventory_to_send.clear();
6354 LOCK(tx_relay->m_bloom_filter_mutex);
6356 for (
const auto& txinfo : vtxinfo) {
6357 const Txid& txid{txinfo.tx->GetHash()};
6358 const Wtxid& wtxid{txinfo.tx->GetWitnessHash()};
6359 const auto inv = peer.m_wtxid_relay ?
6364 if (txinfo.fee < filterrate.GetFee(txinfo.vsize)) {
6367 if (tx_relay->m_bloom_filter) {
6368 if (!tx_relay->m_bloom_filter->IsRelevantAndUpdate(*txinfo.tx))
continue;
6370 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6371 vInv.push_back(inv);
6383 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6386 auto& invs = tx_relay->m_tx_inventory_to_send;
6387 std::vector<CTransactionRef> res;
6389 if (invs.size() == 0)
return res;
6392 if (invs.capacity() > 2 * invs.size()) invs.shrink_to_fit();
6396 res.reserve(txiters.size());
6397 for (
auto txiter : txiters) {
6398 if (txiter->GetFee() < filterrate.GetFee(txiter->GetTxSize())) {
6401 res.push_back(txiter->GetSharedTx());
6404 tx_relay->m_last_inv_sequence = m_mempool.
GetSequence();
6408 LOCK(tx_relay->m_bloom_filter_mutex);
6409 vInv.reserve(std::min<size_t>(
MAX_INV_SZ, vInv.size() + inv_tx.size()));
6410 for (
auto& tx : inv_tx) {
6414 const auto inv = peer.m_wtxid_relay ?
6418 if (tx_relay->m_tx_inventory_known_filter.contains(inv.
hash)) {
6421 if (tx_relay->m_bloom_filter && !tx_relay->m_bloom_filter->IsRelevantAndUpdate(*tx))
continue;
6423 vInv.push_back(inv);
6428 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6436 auto stalling_timeout = m_block_stalling_timeout.load();
6437 if (state.m_stalling_since.count() && state.m_stalling_since < current_time - stalling_timeout) {
6441 LogInfo(
"Peer is stalling block download, %s",
node.DisconnectMsg());
6442 node.fDisconnect =
true;
6446 if (stalling_timeout != new_timeout && m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
6456 if (state.vBlocksInFlight.size() > 0) {
6457 QueuedBlock &queuedBlock = state.vBlocksInFlight.front();
6458 int nOtherPeersWithValidatedDownloads = m_peers_downloading_from - 1;
6460 LogInfo(
"Timeout downloading block %s, %s", queuedBlock.pindex->GetBlockHash().ToString(),
node.DisconnectMsg());
6461 node.fDisconnect =
true;
6466 if (state.fSyncStarted && peer.m_headers_sync_timeout < std::chrono::microseconds::max()) {
6468 if (m_chainman.m_best_header->Time() <=
NodeClock::now() - 24h) {
6469 if (current_time > peer.m_headers_sync_timeout && nSyncStarted == 1 && (m_num_preferred_download_peers - state.fPreferredDownload >= 1)) {
6476 LogInfo(
"Timeout downloading headers, %s",
node.DisconnectMsg());
6477 node.fDisconnect =
true;
6480 LogInfo(
"Timeout downloading headers from noban peer, not %s",
node.DisconnectMsg());
6486 state.fSyncStarted =
false;
6488 peer.m_headers_sync_timeout = 0us;
6494 peer.m_headers_sync_timeout = std::chrono::microseconds::max();
6500 ConsiderEviction(
node, peer, GetTime<std::chrono::seconds>());
6505 std::vector<CInv> vGetData;
6507 std::vector<const CBlockIndex*> vToDownload;
6509 auto get_inflight_budget = [&state]() {
6516 FindNextBlocksToDownload(peer, get_inflight_budget(), vToDownload, staller);
6517 auto historical_blocks{m_chainman.GetHistoricalBlockRange()};
6518 if (historical_blocks && !IsLimitedPeer(peer)) {
6522 TryDownloadingHistoricalBlocks(
6524 get_inflight_budget(),
6525 vToDownload, from_tip, historical_blocks->second);
6528 uint32_t nFetchFlags = GetFetchFlags(peer);
6530 BlockRequested(
node.GetId(), *pindex);
6534 if (state.vBlocksInFlight.empty() && staller != -1) {
6535 if (
State(staller)->m_stalling_since == 0us) {
6536 State(staller)->m_stalling_since = current_time;
6546 LOCK(m_tx_download_mutex);
6547 for (
const GenTxid& gtxid : m_txdownloadman.GetRequestsToSend(
node.GetId(), current_time)) {
6556 if (!vGetData.empty())
6559 MaybeSendFeefilter(
node, peer, current_time);
constexpr CAmount MAX_MONEY
No amount larger than this (in satoshi) is valid.
bool MoneyRange(const CAmount &nValue)
int64_t CAmount
Amount in satoshis (Can be negative)
enum ReadStatus_t ReadStatus
const std::string & BlockFilterTypeName(BlockFilterType filter_type)
Get the human-readable name for a filter type.
BlockFilterIndex * GetBlockFilterIndex(BlockFilterType filter_type)
Get a block filter index by type.
constexpr int CFCHECKPT_INTERVAL
Interval between compact filter checkpoints.
CBlockLocator GetLocator(const CBlockIndex *index)
Get a locator for a block index entry.
int64_t GetBlockProofEquivalentTime(const CBlockIndex &to, const CBlockIndex &from, const CBlockIndex &tip, const Consensus::Params ¶ms)
Return the time it would take to redo the work difference between from and to, assuming the current h...
const CBlockIndex * LastCommonAncestor(const CBlockIndex *pa, const CBlockIndex *pb)
Find the last common ancestor two blocks have.
@ BLOCK_VALID_CHAIN
Outputs do not overspend inputs, no double spends, coinbase output ok, no immature coinbase spends,...
@ BLOCK_VALID_TRANSACTIONS
Only first tx is coinbase, 2 <= coinbase input script length <= 100, transactions valid,...
@ BLOCK_VALID_SCRIPTS
Scripts & signatures ok.
@ BLOCK_VALID_TREE
All parent headers found, difficulty matches, timestamp >= median previous.
@ BLOCK_HAVE_DATA
full block available in blk*.dat
arith_uint256 GetBlockProof(const CBlockIndex &block)
Compute how much work a block index entry corresponds to.
#define Assert(val)
Identity function.
#define Assume(val)
Assume is the identity function.
Stochastic address manager.
void Connected(const CService &addr, NodeSeconds time=Now< NodeSeconds >())
We have successfully connected to this peer.
bool Good(const CService &addr, NodeSeconds time=Now< NodeSeconds >())
Mark an address record as accessible and attempt to move it to addrman's tried table.
bool Add(const std::vector< CAddress > &vAddr, const CNetAddr &source, std::chrono::seconds time_penalty=0s)
Attempt to add one or more addresses to addrman's new table.
void SetServices(const CService &addr, ServiceFlags nServices)
Update an entry's service bits.
bool IsBanned(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
Return whether net_addr is banned.
bool IsDiscouraged(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
Return whether net_addr is discouraged.
void Discourage(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
BlockFilterIndex is used to store and retrieve block filters, hashes, and headers for a range of bloc...
bool LookupFilterRange(int start_height, const CBlockIndex *stop_index, std::vector< BlockFilter > &filters_out) const
Get a range of filters between two heights on a chain.
bool LookupFilterHashRange(int start_height, const CBlockIndex *stop_index, std::vector< uint256 > &hashes_out) const
Get a range of filter hashes between two heights on a chain.
bool LookupFilterHeader(const CBlockIndex *block_index, uint256 &header_out) EXCLUSIVE_LOCKS_REQUIRED(!m_cs_headers_cache)
Get a single filter header by block.
std::vector< CTransactionRef > txn
std::vector< uint16_t > indexes
A CService with information about it as peer.
ServiceFlags nServices
Serialized as uint64_t in V1, and as CompactSize in V2.
static constexpr SerParams V1_NETWORK
NodeSeconds nTime
Always included in serialization. The behavior is unspecified if the value is not representable as ui...
static constexpr SerParams V2_NETWORK
size_t BlockTxCount() const
std::vector< CTransactionRef > vtx
The block chain is a tree shaped structure starting with the genesis block at the root,...
bool IsValid(enum BlockStatus nUpTo) const EXCLUSIVE_LOCKS_REQUIRED(
Check whether this block index entry is valid up to the passed validity level.
CBlockIndex * pprev
pointer to the index of the predecessor of this block
CBlockHeader GetBlockHeader() const
arith_uint256 nChainWork
(memory only) Total amount of work (expected number of hashes) in the chain up to and including this ...
bool HaveNumChainTxs() const
Check whether this block and all previous blocks back to the genesis block or an assumeutxo snapshot ...
uint256 GetBlockHash() const
int64_t GetBlockTime() const
unsigned int nTx
Number of transactions in this block.
CBlockIndex * GetAncestor(int height)
Efficiently find an ancestor of this block.
int nHeight
height of the entry in the chain. The genesis block has height 0
FlatFilePos GetBlockPos() const EXCLUSIVE_LOCKS_REQUIRED(
BloomFilter is a probabilistic filter which SPV clients provide so that we can filter the transaction...
bool IsWithinSizeConstraints() const
True if the size is <= MAX_BLOOM_FILTER_SIZE and the number of hash functions is <= MAX_HASH_FUNCS (c...
An in-memory indexed chain of blocks.
bool Contains(const CBlockIndex &index) const
Efficiently check whether a block is present in this chain.
CBlockIndex * Tip() const
Returns the index entry for the tip of this chain, or nullptr if none.
CBlockIndex * Next(const CBlockIndex &index) const
Find the successor of a block in this chain, or nullptr if the given index is not found or is the tip...
int Height() const
Return the maximal height in the chain.
CChainParams defines various tweakable parameters of a given instance of the Bitcoin system.
const HeadersSyncParams & HeadersSync() const
const Consensus::Params & GetConsensus() const
void NumToOpenAdd(size_t n)
Increment the number of new connections of type ConnectionType::PRIVATE_BROADCAST to be opened by CCo...
size_t NumToOpenSub(size_t n)
Decrement the number of new connections of type ConnectionType::PRIVATE_BROADCAST to be opened by CCo...
bool GetNetworkActive() const
bool GetTryNewOutboundPeer() const
class CConnman::PrivateBroadcast m_private_broadcast
bool ShouldRunInactivityChecks(const CNode &node, NodeClock::time_point now) const
Return true if we should disconnect the peer for failing an inactivity check.
std::vector< CAddress > GetAddresses(CNode &requestor, size_t max_addresses, size_t max_pct)
Return addresses from the per-requestor cache.
void SetTryNewOutboundPeer(bool flag)
void WakeMessageHandler() EXCLUSIVE_LOCKS_REQUIRED(!mutexMsgProc)
bool OutboundTargetReached(bool historicalBlockServingLimit) const EXCLUSIVE_LOCKS_REQUIRED(!m_total_bytes_sent_mutex)
check if the outbound target is reached if param historicalBlockServingLimit is set true,...
void StartExtraBlockRelayPeers()
void ForEachNode(const NodeFn &func) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
CSipHasher GetDeterministicRandomizer(uint64_t id) const
Get a unique deterministic randomizer.
bool EvictTxPeerIfFull(std::optional< NodeId > protect_peer=std::nullopt) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
If we are at capacity for inbound tx-relay peers, attempt to evict one.
uint32_t GetMappedAS(const CNetAddr &addr) const
bool ForNode(NodeId id, std::function< bool(CNode *pnode)> func) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
std::vector< CAddress > GetAddressesUnsafe(size_t max_addresses, size_t max_pct, std::optional< Network > network, bool filtered=true) const
Return randomly selected addresses.
int GetExtraBlockRelayCount() const EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
bool DisconnectNode(std::string_view node) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
Mutex & GetNodesMutex() const LOCK_RETURNED(m_nodes_mutex)
int GetExtraFullOutboundCount() const EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
bool GetUseAddrmanOutgoing() const
bool CheckIncomingNonce(uint64_t nonce) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
Fee rate in satoshis per virtualbyte: CAmount / vB the feerate is represented internally as FeeFrac.
CAmount GetFeePerK() const
Return the fee in satoshis for a vsize of 1000 vbytes.
bool IsMsgCmpctBlk() const
std::string ToString() const
bool IsMsgFilteredBlk() const
bool IsMsgWitnessBlk() const
Used to relay blocks as header + vector<merkle branch> to filtered nodes.
std::vector< std::pair< unsigned int, Txid > > vMatchedTxn
Public only for unit testing and relay testing (not relayed).
bool IsRelayable() const
Whether this address should be relayed to other peers even if we can't reach it ourselves.
static constexpr SerParams V1
bool IsAddrV1Compatible() const
Check if the current object can be serialized in pre-ADDRv2/BIP155 format.
Transport protocol agnostic message container.
Information about a peer.
bool IsFeelerConn() const
bool ExpectServicesFromConn() const
std::atomic< int > nVersion
std::atomic_bool m_has_all_wanted_services
Whether this peer provides all services that we want.
bool IsInboundConn() const
bool HasPermission(NetPermissionFlags permission) const
std::string LogPeer() const
Helper function to log the peer id, optionally including IP address.
bool IsOutboundOrBlockRelayConn() const
bool IsManualConn() const
std::string ConnectionTypeAsString() const
void SetCommonVersion(int greatest_common_version)
std::atomic< bool > m_bip152_highbandwidth_to
std::atomic_bool m_relays_txs
Whether we should relay transactions to this peer.
std::atomic< bool > m_bip152_highbandwidth_from
std::atomic_bool fSuccessfullyConnected
fSuccessfullyConnected is set to true on receiving VERACK from the peer.
bool IsAddrFetchConn() const
uint64_t GetLocalNonce() const
void SetAddrLocal(const CService &addrLocalIn) EXCLUSIVE_LOCKS_REQUIRED(!m_addr_local_mutex)
May not be called more than once.
const NodeClock::time_point m_connected
Unix epoch time at peer connection.
bool IsBlockOnlyConn() const
int GetCommonVersion() const
bool IsFullOutboundConn() const
std::atomic_bool fPauseSend
std::string DisconnectMsg() const
Helper function to log disconnects.
void PongReceived(NodeClock::duration ping_time)
A ping-pong round trip has completed successfully. Update latest and minimum ping durations.
std::atomic_bool m_bloom_filter_loaded
Whether this peer has loaded a bloom filter.
bool IsPrivateBroadcastConn() const
const std::unique_ptr< Transport > m_transport
Transport serializer/deserializer.
const bool m_inbound_onion
Whether this peer is an inbound onion, i.e. connected via our Tor onion service.
std::atomic< std::chrono::seconds > m_last_block_time
UNIX epoch time of the last block received from this peer that we had not yet seen (e....
int AdvertisedVersion() const
Protocol version advertised in our VERSION message.
std::atomic_bool fDisconnect
std::atomic< std::chrono::seconds > m_last_tx_time
UNIX epoch time of the last transaction received from this peer that we had not yet seen (e....
RollingBloomFilter is a probabilistic "keep track of most recently inserted" set.
Simple class for background tasks that should be run periodically or once "after a while".
void scheduleEvery(Function f, std::chrono::milliseconds delta) EXCLUSIVE_LOCKS_REQUIRED(!newTaskMutex)
Repeat f until the scheduler is stopped.
void scheduleFromNow(Function f, std::chrono::milliseconds delta) EXCLUSIVE_LOCKS_REQUIRED(!newTaskMutex)
Call f once after the delta has passed.
A combination of a network address (CNetAddr) and a (TCP) port.
std::string ToStringAddrPort() const
std::vector< unsigned char > GetKey() const
General SipHash-2-4 implementation.
uint64_t Finalize() const
Compute the 64-bit SipHash-2-4 of the data written so far.
CSipHasher & Write(uint64_t data)
Hash a 64-bit integer worth of data.
CTxMemPool stores valid-according-to-the-current-best-chain transactions that may be included in the ...
TxMempoolInfo info_for_relay(const T &id, uint64_t last_sequence) const
Returns info for a transaction if its entry_sequence < last_sequence.
CFeeRate GetMinFee(size_t sizelimit) const
CTransactionRef get(const Txid &hash) const
Return a mempool transaction with a given hash.
size_t DynamicMemoryUsage() const
std::vector< TxMempoolInfo > infoAll() const
bool exists(const Txid &txid) const
uint64_t GetSequence() const EXCLUSIVE_LOCKS_REQUIRED(cs)
std::set< Txid > GetUnbroadcastTxs() const
Returns transactions in unbroadcast set.
unsigned long size() const
void RemoveUnbroadcastTx(const Txid &txid, bool unchecked=false)
Removes a transaction from the unbroadcast set.
std::vector< txiter > ExtractBestByMiningScoreWithTopology(std::vector< Wtxid > &wtxids, size_t n_to_sort) const EXCLUSIVE_LOCKS_REQUIRED(cs)
Look up wtxids in the mempool and (partially) sort by mining score.
virtual void NewPoWValidBlock(const CBlockIndex *pindex, const std::shared_ptr< const CBlock > &block)
Notifies listeners that a block which builds directly on our current tip has been received and connec...
virtual void UpdatedBlockTip(const CBlockIndex *pindexNew, const CBlockIndex *pindexFork, bool fInitialDownload)
Notifies listeners when the block chain tip advances.
virtual void BlockChecked(const std::shared_ptr< const CBlock > &, const BlockValidationState &)
Notifies listeners of a block validation result.
virtual void ActiveTipChange(const CBlockIndex &new_tip, bool is_ibd)
Notifies listeners any time the block chain tip changes, synchronously.
virtual void BlockDisconnected(const std::shared_ptr< const CBlock > &block, const CBlockIndex *pindex)
Notifies listeners of a block being disconnected Provides the block that was disconnected.
virtual void BlockConnected(const kernel::ChainstateRole &role, const std::shared_ptr< const CBlock > &block, const CBlockIndex *pindex)
Notifies listeners of a block being connected.
void ClearBlockIndexCandidates() EXCLUSIVE_LOCKS_REQUIRED(void PopulateBlockIndexCandidates() EXCLUSIVE_LOCKS_REQUIRED(const CBlockIndex * FindForkInGlobalIndex(const CBlockLocator &locator) const EXCLUSIVE_LOCKS_REQUIRED(cs_main)
Populate the candidate set by calling TryAddBlockIndexCandidate on all valid block indices.
Interface for managing multiple Chainstate objects, where each chainstate is associated with chainsta...
bool IsInitialBlockDownload() const noexcept
Check whether we are doing an initial block download (synchronizing from disk or network)
MempoolAcceptResult ProcessTransaction(const CTransactionRef &tx, bool test_accept=false) EXCLUSIVE_LOCKS_REQUIRED(cs_main)
Try to add a transaction to the memory pool.
RecursiveMutex & GetMutex() const LOCK_RETURNED(
Alias for cs_main.
CBlockIndex * ActiveTip() const EXCLUSIVE_LOCKS_REQUIRED(GetMutex())
Chainstate & ActiveChainstate() const
Alternatives to CurrentChainstate() used by older code to query latest chainstate information without...
SnapshotCompletionResult MaybeValidateSnapshot(Chainstate &validated_cs, Chainstate &unvalidated_cs) EXCLUSIVE_LOCKS_REQUIRED(Chainstate & CurrentChainstate() const EXCLUSIVE_LOCKS_REQUIRED(GetMutex())
Try to validate an assumeutxo snapshot by using a validated historical chainstate targeted at the sna...
bool ProcessNewBlock(const std::shared_ptr< const CBlock > &block, bool force_processing, bool min_pow_checked, bool *new_block) LOCKS_EXCLUDED(cs_main)
Process an incoming block.
bool ProcessNewBlockHeaders(std::span< const CBlockHeader > headers, bool min_pow_checked, BlockValidationState &state, const CBlockIndex **ppindex=nullptr) LOCKS_EXCLUDED(cs_main)
Process incoming block headers.
const arith_uint256 & MinimumChainWork() const
CChain & ActiveChain() const EXCLUSIVE_LOCKS_REQUIRED(GetMutex())
void ReportHeadersPresync(int64_t height, int64_t timestamp)
This is used by net_processing to report pre-synchronization progress of headers, as headers are not ...
node::BlockManager m_blockman
A single BlockManager instance is shared across each constructed chainstate to avoid duplicating bloc...
Double ended buffer combining vector and stream-like interfaces.
void ignore(size_t num_ignore)
uint64_t rand64() noexcept
Generate a random 64-bit integer.
const uint256 & ToUint256() const LIFETIMEBOUND
static Mutex g_msgproc_mutex
Mutex for anything that is only accessed via the msg processing thread.
virtual void FinalizeNode(const CNode &node)=0
Handle removal of a peer (clear state)
virtual bool ProcessMessages(CNode &node, std::atomic< bool > &interrupt) EXCLUSIVE_LOCKS_REQUIRED(g_msgproc_mutex)=0
Process protocol messages received from a given node.
virtual bool HasAllDesirableServiceFlags(ServiceFlags services) const =0
Callback to determine whether the given set of service flags are sufficient for a peer to be "relevan...
virtual bool SendMessages(CNode &node) EXCLUSIVE_LOCKS_REQUIRED(g_msgproc_mutex)=0
Send queued protocol messages to a given node.
virtual void InitializeNode(const CNode &node, ServiceFlags our_services)=0
Initialize a peer (setup state)
static bool HasFlag(NetPermissionFlags flags, NetPermissionFlags f)
ReadStatus FillBlock(CBlock &block, const std::vector< CTransactionRef > &vtx_missing, bool segwit_active)
bool IsTxAvailable(size_t index) const
ReadStatus InitData(const CBlockHeaderAndShortTxIDs &cmpctblock, const std::vector< std::pair< Wtxid, CTransactionRef > > &extra_txn)
virtual util::Expected< void, std::string > FetchBlock(NodeId peer_id, const CBlockIndex &block_index)=0
Attempt to manually fetch block from a given peer.
virtual ServiceFlags GetDesirableServiceFlags(ServiceFlags services) const =0
Gets the set of service flags which are "desirable" for a given peer.
virtual void StartScheduledTasks(CScheduler &scheduler)=0
Begin running background tasks, should only be called once.
virtual std::vector< node::TxOrphanage::OrphanInfo > GetOrphanTransactions()=0
static std::unique_ptr< PeerManager > make(CConnman &connman, AddrMan &addrman, BanMan *banman, ChainstateManager &chainman, CTxMemPool &pool, node::Warnings &warnings, Options opts)
virtual void UnitTestMisbehaving(NodeId peer_id)=0
virtual bool GetNodeStateStats(NodeId nodeid, CNodeStateStats &stats) const =0
Get statistics from node state.
virtual void UpdateLastBlockAnnounceTime(NodeId node, int64_t time_in_seconds)=0
This function is used for testing the stale tip eviction logic, see denialofservice_tests....
virtual void CheckForStaleTipAndEvictPeers()=0
Evict extra outbound peers.
Store a list of transactions to be broadcast privately.
@ QueueFull
Rejected: the queue is already at MAX_TRANSACTIONS.
@ AlreadyPresent
The transaction was already present with send attempts remaining; no change.
@ Added
The transaction was newly added or reset after exhausting its send attempts.
static constexpr size_t MAX_TRANSACTIONS
Maximum number of transactions tracked simultaneously.
I randrange(I range) noexcept
Generate a random integer in the range [0..range), with range > 0.
bool Contains(Network net) const EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
std::string GetDebugMessage() const
std::string ToString() const
256-bit unsigned big integer.
constexpr bool IsNull() const
std::string ToString() const
CBlockIndex * LookupBlockIndex(const uint256 &hash) EXCLUSIVE_LOCKS_REQUIRED(cs_main)
CBlockFileInfo *GetBlockFileInfo(size_t n) EXCLUSIVE_LOCKS_REQUIRED(bool WriteBlockUndo(const CBlockUndo &blockundo, BlockValidationState &state, CBlockIndex &block) EXCLUSIVE_LOCKS_REQUIRED(FlatFilePos WriteBlock(const CBlock &block, int nHeight) EXCLUSIVE_LOCKS_REQUIRED(void UpdateBlockInfo(const CBlock &block, unsigned int nHeight, const FlatFilePos &pos) EXCLUSIVE_LOCKS_REQUIRED(bool IsPruneMode() const
Get block file info entry for one block file.
bool LoadingBlocks() const
ReadRawBlockResult ReadRawBlock(const FlatFilePos &pos, std::optional< std::pair< size_t, size_t > > block_part=std::nullopt) const
bool ReadBlock(CBlock &block, const FlatFilePos &pos, const std::optional< uint256 > &expected_hash) const
Functions for disk access for blocks.
Class responsible for deciding what transactions to request and, once downloaded, whether and how to ...
Manages warning messages within a node.
std::string ToString() const
const uint256 & ToUint256() const LIFETIMEBOUND
The util::Expected class provides a standard way for low-level functions to return either error value...
A token bucket rate limiter.
bool decrement(double n=1.0, double floor=0.0)
Consume n tokens.
void increment(const time_point &now)
Refill tokens based on elapsed time since last call.
double value() const
Current token balance.
The util::Unexpected class represents an unexpected value stored in util::Expected.
std::string TransportTypeAsString(TransportProtocolType transport_type)
Convert TransportProtocolType enum to a string value.
@ BLOCK_HEADER_LOW_WORK
the block header may be on a too-little-work chain
@ BLOCK_INVALID_HEADER
invalid proof of work or time too old
@ BLOCK_CACHED_INVALID
this block was cached as being invalid and we didn't store the reason why
@ BLOCK_CONSENSUS
invalid by consensus rules (excluding any below reasons)
@ BLOCK_MISSING_PREV
We don't have the previous block the checked one is built on.
@ BLOCK_INVALID_PREV
A block this one builds on is invalid.
@ BLOCK_MUTATED
the block's data didn't match the data committed to by the PoW
@ BLOCK_TIME_FUTURE
block timestamp was > 2 hours in the future (or our clock is bad)
@ BLOCK_RESULT_UNSET
initial value. Block has not yet been rejected
@ TX_MISSING_INPUTS
transaction was missing some of its inputs
@ TX_UNKNOWN
transaction was not validated because package failed
@ TX_NO_MEMPOOL
this node does not have a mempool so can't validate the transaction
@ TX_RESULT_UNSET
initial value. Tx has not yet been rejected
static size_t RecursiveDynamicUsage(const CScript &script)
RecursiveMutex cs_main
Mutex to guard access to validation specific variables, such as reading or changing the chainstate.
bool DeploymentActiveAfter(const CBlockIndex *pindexPrev, const Consensus::Params ¶ms, Consensus::BuriedDeployment dep, VersionBitsCache &versionbitscache)
Determine if a deployment is active for the next block.
bool DeploymentActiveAt(const CBlockIndex &index, const Consensus::Params ¶ms, Consensus::BuriedDeployment dep, VersionBitsCache &versionbitscache)
Determine if a deployment is active for this block.
is a home for simple enum and struct type definitions that can be used internally by functions in the...
#define LogDebug(category,...)
CSerializedNetMsg Make(std::string msg_type, Args &&... args)
constexpr const char * FILTERCLEAR
The filterclear message tells the receiving peer to remove a previously-set bloom filter.
constexpr const char * FEEFILTER
The feefilter message tells the receiving peer not to inv us any txs which do not meet the specified ...
constexpr const char * SENDHEADERS
Indicates that a node prefers to receive new block announcements via a "headers" message rather than ...
constexpr const char * GETBLOCKS
The getblocks message requests an inv message that provides block header hashes starting from a parti...
constexpr const char * HEADERS
The headers message sends one or more block headers to a node which previously requested certain head...
constexpr const char * ADDR
The addr (IP address) message relays connection information for peers on the network.
constexpr const char * GETBLOCKTXN
Contains a BlockTransactionsRequest Peer should respond with "blocktxn" message.
constexpr const char * CMPCTBLOCK
Contains a CBlockHeaderAndShortTxIDs object - providing a header and list of "short txids".
constexpr const char * CFCHECKPT
cfcheckpt is a response to a getcfcheckpt request containing a vector of evenly spaced filter headers...
constexpr const char * SENDADDRV2
The sendaddrv2 message signals support for receiving ADDRV2 messages (BIP155).
constexpr const char * GETADDR
The getaddr message requests an addr message from the receiving node, preferably one with lots of IP ...
constexpr const char * GETCFILTERS
getcfilters requests compact filters for a range of blocks.
constexpr const char * PONG
The pong message replies to a ping message, proving to the pinging node that the ponging node is stil...
constexpr const char * BLOCKTXN
Contains a BlockTransactions.
constexpr const char * CFHEADERS
cfheaders is a response to a getcfheaders request containing a filter header and a vector of filter h...
constexpr const char * PING
The ping message is sent periodically to help confirm that the receiving peer is still connected.
constexpr const char * FILTERLOAD
The filterload message tells the receiving peer to filter all relayed transactions and requested merk...
constexpr const char * SENDTXRCNCL
Contains a 4-byte version number and an 8-byte salt.
constexpr const char * ADDRV2
The addrv2 message relays connection information for peers on the network just like the addr message,...
constexpr const char * VERACK
The verack message acknowledges a previously-received version message, informing the connecting node ...
constexpr const char * GETHEADERS
The getheaders message requests a headers message that provides block headers starting from a particu...
constexpr const char * FILTERADD
The filteradd message tells the receiving peer to add a single element to a previously-set bloom filt...
constexpr const char * CFILTER
cfilter is a response to a getcfilters request containing a single compact filter.
constexpr const char * FEATURE
BIP 434 Peer feature negotiation.
constexpr const char * GETDATA
The getdata message requests one or more data objects from another node.
constexpr const char * SENDCMPCT
Contains a 1-byte bool and 8-byte LE version number.
constexpr const char * GETCFCHECKPT
getcfcheckpt requests evenly spaced compact filter headers, enabling parallelized download and valida...
constexpr const char * INV
The inv message (inventory message) transmits one or more inventories of objects known to the transmi...
constexpr const char * TX
The tx message transmits a single transaction.
constexpr const char * MEMPOOL
The mempool message requests the TXIDs of transactions that the receiving node has verified as valid ...
constexpr const char * NOTFOUND
The notfound message is a reply to a getdata message which requested an object the receiving node doe...
constexpr const char * MERKLEBLOCK
The merkleblock message is a reply to a getdata message which requested a block using the inventory t...
constexpr const char * WTXIDRELAY
Indicates that a node prefers to relay transactions via wtxid, rather than txid.
constexpr const char * BLOCK
The block message transmits a single serialized block.
constexpr const char * GETCFHEADERS
getcfheaders requests a compact filter header and the filter hashes for a range of blocks,...
constexpr const char * VERSION
The version message provides information about the transmitting node to the receiving node at the beg...
constexpr int32_t MAX_PEER_TX_ANNOUNCEMENTS
Maximum number of transactions to consider for requesting, per peer.
""_hex is a compile-time user-defined literal returning a std::array<std::byte>, equivalent to ParseH...
bool ShouldDebugLog(Category category)
Return whether messages with specified category should be debug logged.
std::string ToString(const T &t)
Locale-independent version of std::to_string.
std::string strSubVersion
Subversion as sent to the P2P network in version messages.
std::optional< CService > GetLocalAddrForPeer(CNode &node)
Returns a local address that we should advertise to this peer.
std::function< void(const CAddress &addr, const std::string &msg_type, std::span< const unsigned char > data, bool is_incoming)> CaptureMessage
Defaults to CaptureMessageToFile(), but can be overridden by unit tests.
bool SeenLocal(const CService &addr)
vote for a local address
constexpr unsigned int MAX_SUBVERSION_LENGTH
Maximum length of the user agent string in version message.
constexpr std::chrono::minutes TIMEOUT_INTERVAL
Time after which to disconnect, after waiting for a ping response (or inactivity).
static constexpr auto HEADERS_RESPONSE_TIME
How long to wait for a peer to respond to a getheaders request.
static constexpr size_t MAX_ADDR_TO_SEND
The maximum number of address records permitted in an ADDR message.
static constexpr auto INVENTORY_BUCKET_CHECK_DELAY
Delay between checking inventory bucket and backlog.
static constexpr size_t MAX_ADDR_PROCESSING_TOKEN_BUCKET
The soft limit of the address processing token bucket (the regular MAX_ADDR_RATE_PER_SECOND based inc...
static constexpr auto INVENTORY_BUCKET_BACKLOG_HEARTBEAT
Delay between inventory bucket backlog heartbeat log entries.
TRACEPOINT_SEMAPHORE(net, inbound_message)
static const int MAX_BLOCKS_IN_TRANSIT_PER_PEER
Number of blocks that can be requested at any given time from a single peer.
static constexpr auto BLOCK_STALLING_TIMEOUT_DEFAULT
Default time during which a peer must stall block download progress before being disconnected.
static constexpr auto AVG_FEEFILTER_BROADCAST_INTERVAL
Average delay between feefilter broadcasts in seconds.
static constexpr auto EXTRA_PEER_CHECK_INTERVAL
How frequently to check for extra outbound peers and disconnect.
static const unsigned int BLOCK_DOWNLOAD_WINDOW
Size of the "block download window": how far ahead of our current height do we fetch?...
static constexpr int STALE_RELAY_AGE_LIMIT
Age after which a stale block will no longer be served if requested as protection against fingerprint...
static constexpr int HISTORICAL_BLOCK_AGE
Age after which a block is considered historical for purposes of rate limiting block relay.
static constexpr auto ROTATE_ADDR_RELAY_DEST_INTERVAL
Delay between rotating the peers we relay a particular address to.
static constexpr auto MINIMUM_CONNECT_TIME
Minimum time an outbound-peer-eviction candidate must be connected for, in order to evict.
static constexpr auto CHAIN_SYNC_TIMEOUT
Timeout for (unprotected) outbound peers to sync to our chainwork.
static constexpr auto OUTBOUND_INVENTORY_BROADCAST_INTERVAL
Average delay between trickled inventory transmissions for outbound peers.
static const unsigned int NODE_NETWORK_LIMITED_MIN_BLOCKS
Minimum blocks required to signal NODE_NETWORK_LIMITED.
static constexpr auto AVG_LOCAL_ADDRESS_BROADCAST_INTERVAL
Average delay between local address broadcasts.
static const int MAX_BLOCKTXN_DEPTH
Maximum depth of blocks we're willing to respond to GETBLOCKTXN requests for.
static constexpr size_t INVENTORY_BUCKET_BACKLOG_CAPACITY
Empty backlog target capacity.
static constexpr int32_t MAX_OUTBOUND_PEERS_TO_PROTECT_FROM_DISCONNECT
Protect at least this many outbound peers from disconnection due to slow/ behind headers chain.
static constexpr auto INBOUND_INVENTORY_BROADCAST_INTERVAL
Average delay between trickled inventory transmissions for inbound peers.
static constexpr size_t NUM_PRIVATE_BROADCAST_PER_TX
For private broadcast, send a transaction to this many peers.
static constexpr auto MAX_FEEFILTER_CHANGE_DELAY
Maximum feefilter broadcast delay after significant change.
static constexpr uint32_t MAX_GETCFILTERS_SIZE
Maximum number of compact filters that may be requested with one getcfilters.
static constexpr double OUTBOUND_INVENTORY_BUCKET_MULTIPLIER
Multiplier for the inventory bucket rate for outbounds.
static constexpr auto HEADERS_DOWNLOAD_TIMEOUT_BASE
Headers download timeout.
static const unsigned int MAX_GETDATA_SZ
Limit to avoid sending big packets.
static constexpr double BLOCK_DOWNLOAD_TIMEOUT_BASE
Block download timeout base, expressed in multiples of the block interval (i.e.
static constexpr auto PRIVATE_BROADCAST_MAX_CONNECTION_LIFETIME
Private broadcast connections must complete within this time.
static constexpr auto STALE_CHECK_INTERVAL
How frequently to check for stale tips.
static constexpr auto AVG_ADDRESS_BROADCAST_INTERVAL
Average delay between peer address broadcasts.
static const unsigned int MAX_LOCATOR_SZ
The maximum number of entries in a locator.
static constexpr double BLOCK_DOWNLOAD_TIMEOUT_PER_PEER
Additional block download timeout per parallel downloading peer (i.e.
static constexpr double MAX_ADDR_RATE_PER_SECOND
The maximum rate of address records we're willing to process on average.
static constexpr auto PING_INTERVAL
Time between pings automatically sent out for latency probing and keepalive.
static constexpr size_t INVENTORY_BUCKET_BACKLOG_HEARTBEAT_MIN
Minimum backlog to trigger heartbeat log entries.
static const int MAX_CMPCTBLOCK_DEPTH
Maximum depth of blocks we're willing to serve as compact blocks to peers when requested.
static const unsigned int MAX_BLOCKS_TO_ANNOUNCE
Maximum number of headers to announce when relaying blocks with headers message.
static const unsigned int NODE_NETWORK_LIMITED_ALLOW_CONN_BLOCKS
Window, in blocks, for connecting to NODE_NETWORK_LIMITED peers.
static constexpr uint32_t MAX_GETCFHEADERS_SIZE
Maximum number of cf hashes that may be requested with one getcfheaders.
static constexpr auto BLOCK_STALLING_TIMEOUT_MAX
Maximum timeout for stalling block download.
static constexpr auto HEADERS_DOWNLOAD_TIMEOUT_PER_HEADER
static constexpr uint64_t RANDOMIZER_ID_ADDRESS_RELAY
SHA256("main address relay")[0:8].
static constexpr size_t MAX_PCT_ADDR_TO_SEND
the maximum percentage of addresses from our addrman to return in response to a getaddr message.
static const unsigned int MAX_INV_SZ
The maximum number of entries in an 'inv' protocol message.
constexpr unsigned int MAX_CMPCTBLOCKS_INFLIGHT_PER_BLOCK
Maximum number of outstanding CMPCTBLOCK requests for the same block.
constexpr uint64_t CMPCTBLOCKS_VERSION
The compactblocks version we support.
ReachableNets g_reachable_nets
bool IsProxy(const CNetAddr &addr)
constexpr unsigned int DEFAULT_MIN_RELAY_TX_FEE
Default for -minrelaytxfee, minimum relay fee for transactions.
constexpr TransactionSerParams TX_NO_WITNESS
constexpr TransactionSerParams TX_WITH_WITNESS
std::shared_ptr< const CTransaction > CTransactionRef
GenTxid ToGenTxid(const CInv &inv)
Convert a TX/WITNESS_TX/WTX CInv to a GenTxid.
constexpr size_t MAX_FEATUREDATA_LENGTH
constexpr uint32_t MSG_WITNESS_FLAG
getdata message type flags
@ MSG_WTX
Defined in BIP 339.
@ MSG_CMPCT_BLOCK
Defined in BIP152.
@ MSG_WITNESS_BLOCK
Defined in BIP144.
ServiceFlags
nServices flags
constexpr size_t MAX_FEATUREID_LENGTH
static bool MayHaveUsefulAddressDB(ServiceFlags services)
Checks if a peer with the given service flags may be capable of having a robust address-storage DB.
constexpr int MIN_PEER_PROTO_VERSION
disconnect from peers older than this proto version
constexpr int SHORT_IDS_BLOCKS_VERSION
short-id-based block download starts with this version
constexpr int BIP0031_VERSION
BIP 0031, pong message, is enabled for all versions AFTER this one.
constexpr int FEEFILTER_VERSION
"feefilter" tells peers to filter invs to you by fee starts with this version
constexpr int WTXID_RELAY_VERSION
"wtxidrelay" message type for wtxid-based relay starts with this version
constexpr int INVALID_CB_NO_BAN_VERSION
not banning for invalid compact blocks starts with this version
constexpr int FEATURE_VERSION
"feature" message type for feature negotiation starts with this version
constexpr int SENDHEADERS_VERSION
"sendheaders" message type and announcing blocks with headers starts with this version
constexpr unsigned int MAX_SCRIPT_ELEMENT_SIZE
#define LIMITED_VECTOR(obj, n)
#define LIMITED_STRING(obj, n)
uint64_t ReadCompactSize(Stream &is, bool range_check=true)
Decode a CompactSize-encoded variable-length integer.
constexpr auto MakeUCharSpan(const V &v) -> decltype(UCharSpanCast(std::span{v}))
Like the std::span constructor, but for (const) unsigned char member types only.
Describes a place in the block chain to another node such that if the other node doesn't have the sam...
std::vector< uint256 > vHave
NodeClock::duration m_ping_wait
std::vector< int > vHeightInFlight
CAmount m_fee_filter_received
std::chrono::seconds time_offset
bool m_addr_relay_enabled
uint64_t m_addr_rate_limited
uint64_t m_addr_processed
ServiceFlags their_services
Parameters that influence chain consensus.
int64_t nPowTargetSpacing
std::chrono::seconds PowTargetSpacing() const
Validation result for a transaction evaluated by MemPoolAccept (single or package).
const ResultType m_result_type
Result type.
const TxValidationState m_state
Contains information about why the transaction failed.
@ DIFFERENT_WITNESS
Valid, transaction was already in the mempool.
@ INVALID
Fully validated, valid.
const std::list< CTransactionRef > m_replaced_transactions
Mempool transactions replaced by the tx.
Version of the system clock that is mockable in the context of tests (via FakeNodeClock or SetMockTim...
static time_point now() noexcept
Return current system time or mocked time, if set.
std::chrono::time_point< NodeClock > time_point
static constexpr time_point epoch
Validation result for package mempool acceptance.
PackageValidationState m_state
std::map< Wtxid, MempoolAcceptResult > m_tx_results
Map from wtxid to finished MempoolAcceptResults.
std::chrono::seconds median_outbound_time_offset
Information about chainstate that notifications are sent from.
bool historical
Whether this is a historical chainstate downloading old blocks to validate an assumeutxo snapshot,...
CFeeRate min_relay_feerate
A fee rate smaller than this is considered zero fee (for relaying, mining and transaction creation)
std::vector< NodeId > m_senders
std::string ToString() const
#define AssertLockNotHeld(cs)
#define WITH_LOCK(cs, code)
Run code while locking a mutex.
COutPoint ProcessBlock(const NodeContext &node, const std::shared_ptr< CBlock > &block)
Returns the generated coin (or Null if the block was invalid).
#define EXCLUSIVE_LOCKS_REQUIRED(...)
#define LOCKS_EXCLUDED(...)
#define ACQUIRED_BEFORE(...)
#define TRACEPOINT(context,...)
consteval auto _(util::TranslatedLiteral str)
ReconciliationRegisterResult
constexpr uint32_t TXRECONCILIATION_VERSION
Supported transaction reconciliation protocol version.
std::string SanitizeString(std::string_view str, int rule)
Remove unsafe chars.
int64_t GetTime()
DEPRECATED Use either ClockType::now() or Now<TimePointType>() if a cast is needed.
constexpr int64_t count_seconds(std::chrono::seconds t)
std::chrono::time_point< NodeClock, std::chrono::seconds > NodeSeconds
PackageMempoolAcceptResult ProcessNewPackage(Chainstate &active_chainstate, CTxMemPool &pool, const Package &package, bool test_accept, const std::optional< CFeeRate > &client_maxfeerate)
Validate (and maybe submit) a package to the mempool.
bool IsBlockMutated(const CBlock &block, bool check_witness_root)
Check if a block has been mutated (with respect to its merkle root and witness commitments).
bool HasValidProofOfWork(std::span< const CBlockHeader > headers, const Consensus::Params &consensusParams)
Check that the proof of work on each blockheader matches the value in nBits.
arith_uint256 CalculateClaimedHeadersWork(std::span< const CBlockHeader > headers)
Return the sum of the claimed work on a given set of headers.
@ UNVALIDATED
Blocks after an assumeutxo snapshot have been validated but the snapshot itself has not been validate...
constexpr unsigned int MIN_BLOCKS_TO_KEEP
Block files containing a block-height within MIN_BLOCKS_TO_KEEP of ActiveChain().Tip() will not be pr...