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, m_rng, 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);
3034 m_chainparams.
HeadersSync(), chain_start_header, minimum_chain_work));
3039 (void)IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3053bool PeerManagerImpl::IsAncestorOfBestHeaderOrTip(
const CBlockIndex* header)
3055 if (header ==
nullptr) {
3057 }
else if (m_chainman.m_best_header !=
nullptr && header == m_chainman.m_best_header->GetAncestor(header->
nHeight)) {
3065bool PeerManagerImpl::MaybeSendGetHeaders(
CNode& pfrom,
const CBlockLocator& locator, Peer& peer)
3073 peer.m_last_getheaders_timestamp = current_time;
3084void PeerManagerImpl::HeadersDirectFetchBlocks(
CNode& pfrom,
const Peer& peer,
const CBlockIndex& last_header)
3087 CNodeState *nodestate =
State(pfrom.
GetId());
3090 std::vector<const CBlockIndex*> vToFetch;
3098 vToFetch.push_back(pindexWalk);
3100 pindexWalk = pindexWalk->
pprev;
3112 std::vector<CInv> vGetData;
3114 for (
const CBlockIndex* pindex : vToFetch | std::views::reverse) {
3119 uint32_t nFetchFlags = GetFetchFlags(peer);
3121 BlockRequested(pfrom.
GetId(), *pindex);
3125 if (vGetData.size() > 1) {
3130 if (vGetData.size() > 0) {
3131 if (!m_opts.ignore_incoming_txs &&
3132 nodestate->m_provides_cmpctblocks &&
3133 vGetData.size() == 1 &&
3134 mapBlocksInFlight.size() == 1 &&
3150void PeerManagerImpl::UpdatePeerStateForReceivedHeaders(
CNode& pfrom,
3151 const CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
3154 CNodeState *nodestate =
State(pfrom.
GetId());
3163 nodestate->m_last_block_announcement =
GetTime();
3171 if (nodestate->pindexBestKnownBlock && nodestate->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
3193 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) {
3195 nodestate->m_chain_sync.m_protect =
true;
3196 ++m_outbound_peers_with_protect_from_disconnect;
3201void PeerManagerImpl::ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
3202 std::vector<CBlockHeader>&&
headers,
3203 bool via_compact_block)
3205 size_t nCount =
headers.size();
3212 LOCK(peer.m_headers_sync_mutex);
3213 if (peer.m_headers_sync) {
3214 peer.m_headers_sync.reset(
nullptr);
3215 LOCK(m_headers_presync_mutex);
3216 m_headers_presync_stats.erase(pfrom.
GetId());
3220 peer.m_last_getheaders_timestamp = {};
3228 if (!CheckHeadersPoW(
headers, peer)) {
3243 bool already_validated_work =
false;
3246 bool have_headers_sync =
false;
3248 LOCK(peer.m_headers_sync_mutex);
3250 already_validated_work = IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3266 have_headers_sync = !!peer.m_headers_sync;
3271 bool headers_connect_blockindex{chain_start_header !=
nullptr};
3273 if (!headers_connect_blockindex) {
3277 HandleUnconnectingHeaders(pfrom, peer,
headers);
3284 peer.m_last_getheaders_timestamp = {};
3294 already_validated_work = already_validated_work || IsAncestorOfBestHeaderOrTip(last_received_header);
3301 already_validated_work =
true;
3307 if (!already_validated_work && TryLowWorkHeadersSync(peer, pfrom,
3308 *chain_start_header,
headers)) {
3320 bool received_new_header{last_received_header ==
nullptr};
3326 state, &pindexLast)};
3332 "If this happens with all peers, consider database corruption (that -reindex may fix) "
3333 "or a potential consensus incompatibility.",
3336 MaybePunishNodeForBlock(pfrom.
GetId(), state, via_compact_block,
"invalid header received");
3342 if (processed && received_new_header) {
3343 LogBlockHeader(*pindexLast, pfrom,
false);
3347 if (nCount == m_opts.max_headers_result && !have_headers_sync) {
3349 if (MaybeSendGetHeaders(pfrom,
GetLocator(pindexLast), peer)) {
3354 UpdatePeerStateForReceivedHeaders(pfrom, *pindexLast, received_new_header, nCount == m_opts.max_headers_result);
3357 HeadersDirectFetchBlocks(pfrom, peer, *pindexLast);
3363 bool first_time_failure)
3369 PeerRef peer{GetPeerRef(nodeid)};
3372 ptx->GetHash().ToString(),
3373 ptx->GetWitnessHash().ToString(),
3377 const auto& [add_extra_compact_tx, unique_parents, package_to_validate] = m_txdownloadman.MempoolRejectedTx(ptx, state, nodeid, first_time_failure);
3380 AddToCompactExtraTransactions(ptx);
3382 for (
const Txid& parent_txid : unique_parents) {
3383 if (peer) AddKnownTx(*peer, parent_txid.ToUint256());
3386 return package_to_validate;
3389void PeerManagerImpl::ProcessValidTx(
NodeId nodeid,
const CTransactionRef& tx,
const std::list<CTransactionRef>& replaced_transactions)
3395 m_txdownloadman.MempoolAcceptedTx(tx);
3399 tx->GetHash().ToString(),
3400 tx->GetWitnessHash().ToString(),
3403 InitiateTxBroadcastToAll(tx->GetWitnessHash());
3406 AddToCompactExtraTransactions(removedTx);
3416 const auto&
package = package_to_validate.m_txns;
3417 const auto& senders = package_to_validate.
m_senders;
3420 m_txdownloadman.MempoolRejectedPackage(package);
3423 if (!
Assume(package.size() == 2))
return;
3427 auto package_iter = package.rbegin();
3428 auto senders_iter = senders.rbegin();
3429 while (package_iter != package.rend()) {
3430 const auto& tx = *package_iter;
3431 const NodeId nodeid = *senders_iter;
3432 const auto it_result{package_result.
m_tx_results.find(tx->GetWitnessHash())};
3436 const auto& tx_result = it_result->second;
3437 switch (tx_result.m_result_type) {
3440 ProcessValidTx(nodeid, tx, tx_result.m_replaced_transactions);
3450 ProcessInvalidTx(nodeid, tx, tx_result.m_state,
false);
3468bool PeerManagerImpl::ProcessOrphanTx(Peer& peer)
3473 while (
CTransactionRef porphanTx = m_txdownloadman.GetTxToReconsider(peer.m_id)) {
3476 const Txid& orphanHash = porphanTx->GetHash();
3477 const Wtxid& orphan_wtxid = porphanTx->GetWitnessHash();
3494 ProcessInvalidTx(peer.m_id, porphanTx, state,
false);
3503bool PeerManagerImpl::PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
3505 const uint256& stop_hash, uint32_t max_height_diff,
3509 const bool supported_filter_type =
3512 if (!supported_filter_type) {
3514 static_cast<uint8_t
>(filter_type),
node.DisconnectMsg());
3515 node.fDisconnect =
true;
3524 if (!stop_index || !BlockRequestAllowed(*stop_index)) {
3527 node.fDisconnect =
true;
3532 uint32_t stop_height = stop_index->
nHeight;
3533 if (start_height > stop_height) {
3535 "start height %d and stop height %d, %s",
3536 start_height, stop_height,
node.DisconnectMsg());
3537 node.fDisconnect =
true;
3540 if (stop_height - start_height >= max_height_diff) {
3542 stop_height - start_height + 1, max_height_diff,
node.DisconnectMsg());
3543 node.fDisconnect =
true;
3548 if (!filter_index) {
3558 uint8_t filter_type_ser;
3559 uint32_t start_height;
3562 vRecv >> filter_type_ser >> start_height >> stop_hash;
3568 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3573 std::vector<BlockFilter> filters;
3575 LogDebug(
BCLog::NET,
"Failed to find block filter in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3580 for (
const auto& filter : filters) {
3587 uint8_t filter_type_ser;
3588 uint32_t start_height;
3591 vRecv >> filter_type_ser >> start_height >> stop_hash;
3597 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3603 if (start_height > 0) {
3605 stop_index->
GetAncestor(
static_cast<int>(start_height - 1));
3607 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3613 std::vector<uint256> filter_hashes;
3615 LogDebug(
BCLog::NET,
"Failed to find block filter hashes in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3629 uint8_t filter_type_ser;
3632 vRecv >> filter_type_ser >> stop_hash;
3638 if (!PrepareBlockFilterRequest(
node, peer, filter_type, 0, stop_hash,
3639 std::numeric_limits<uint32_t>::max(),
3640 stop_index, filter_index)) {
3648 for (
int i =
headers.size() - 1; i >= 0; i--) {
3653 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3667 bool new_block{
false};
3668 m_chainman.
ProcessNewBlock(block, force_processing, min_pow_checked, &new_block);
3670 node.m_last_block_time = GetTime<std::chrono::seconds>();
3675 RemoveBlockRequest(block->GetHash(), std::nullopt);
3678 mapBlockSource.erase(block->GetHash());
3682void PeerManagerImpl::ProcessCompactBlockTxns(
CNode& pfrom, Peer& peer,
const BlockTransactions& block_transactions)
3684 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
3685 bool fBlockRead{
false};
3689 auto range_flight = mapBlocksInFlight.equal_range(block_transactions.
blockhash);
3690 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
3691 bool requested_block_from_this_peer{
false};
3694 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
3696 while (range_flight.first != range_flight.second) {
3697 auto [node_id, block_it] = range_flight.first->second;
3698 if (node_id == pfrom.
GetId() && block_it->partialBlock) {
3699 requested_block_from_this_peer =
true;
3702 range_flight.first++;
3705 if (!requested_block_from_this_peer) {
3717 Misbehaving(peer,
"previous compact block reconstruction attempt failed");
3728 Misbehaving(peer,
"invalid compact block/non-matching block transactions");
3731 if (first_in_flight) {
3736 std::vector<CInv> invs;
3741 LogDebug(
BCLog::NET,
"Peer %d sent us a compact block but it failed to reconstruct, waiting on first download to complete\n", pfrom.
GetId());
3754 mapBlockSource.emplace(block_transactions.
blockhash, std::make_pair(pfrom.
GetId(),
false));
3769void PeerManagerImpl::LogBlockHeader(
const CBlockIndex& index,
const CNode& peer,
bool via_compact_block) {
3781 "Saw new %sheader hash=%s height=%d %s",
3782 via_compact_block ?
"cmpctblock " :
"",
3794void PeerManagerImpl::PushPrivateBroadcastTx(
CNode&
node)
3798 const auto opt_tx{m_tx_for_private_broadcast.PickTxForSend(
node.GetId(),
CService{node.addr})};
3801 node.fDisconnect =
true;
3807 tx->GetHash().ToString(), tx->HasWitness() ?
strprintf(
", wtxid=%s", tx->GetWitnessHash().ToString()) :
"",
3813void PeerManagerImpl::ProcessMessage(Peer& peer,
CNode& pfrom,
const std::string& msg_type,
DataStream& vRecv,
3815 const std::atomic<bool>& interruptMsgProc)
3830 uint64_t nNonce = 1;
3833 std::string cleanSubVer;
3834 int starting_height = -1;
3837 vRecv >> nVersion >> Using<CustomUintFormatter<8>>(nServices) >> nTime;
3852 LogDebug(
BCLog::NET,
"peer does not offer the expected services (%08x offered, %08x expected), %s",
3854 GetDesirableServiceFlags(nServices),
3867 if (!vRecv.
empty()) {
3875 if (!vRecv.
empty()) {
3876 std::string strSubVer;
3880 if (!vRecv.
empty()) {
3881 vRecv >> starting_height;
3901 PushNodeVersion(pfrom, peer);
3905 const int greatest_common_version = std::min(nVersion, pfrom.
AdvertisedVersion());
3910 peer.m_their_services = nServices;
3914 pfrom.cleanSubVer = cleanSubVer;
3925 (fRelay || (peer.m_our_services &
NODE_BLOOM))) {
3926 auto*
const tx_relay = peer.SetTxRelay();
3928 LOCK(tx_relay->m_bloom_filter_mutex);
3929 tx_relay->m_relay_txs = fRelay;
3935 LogDebug(
BCLog::NET,
"receive version message: %s: version %d, blocks=%d, us=%s, txrelay=%d, %s%s",
3936 cleanSubVer.empty() ?
"<no user agent>" : cleanSubVer, pfrom.
nVersion,
3938 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
3956 if (greatest_common_version >= 70016) {
3971 const auto* tx_relay = peer.GetTxRelay();
3972 if (tx_relay &&
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs) &&
3974 const uint64_t recon_salt = m_txreconciliation->PreRegisterPeer(pfrom.
GetId());
3987 if (MaybeDisconnectForTxRelayCapacity(pfrom, msg_type, pfrom.
GetId()))
return;
3995 m_num_preferred_download_peers += state->fPreferredDownload;
4001 bool send_getaddr{
false};
4003 send_getaddr = SetupAddressRelay(pfrom, peer);
4013 peer.m_getaddr_sent =
true;
4037 peer.m_time_offset =
NodeSeconds{std::chrono::seconds{nTime}} - Now<NodeSeconds>();
4041 m_outbound_time_offsets.Add(peer.m_time_offset);
4042 m_outbound_time_offsets.WarnIfOutOfSync();
4046 if (greatest_common_version <= 70012) {
4047 constexpr auto finalAlert{
"60010000000000000000000000ffffff7f00000000ffffff7ffeffff7f01ffffff7f00000000ffffff7f00ffffff7f002f555247454e543a20416c657274206b657920636f6d70726f6d697365642c2075706772616465207265717569726564004630440220653febd6410f470f6bae11cad19c48413becb1ac2c17f908fd0fd53bdc3abd5202206d0e9c96fe88d4a0f01ed9dedae2b6f9e00da94cad0fecaae66ecf689bf71b50"_hex};
4048 MakeAndPushMessage(pfrom,
"alert", finalAlert);
4071 auto new_peer_msg = [&]() {
4073 return strprintf(
"New %s peer connected: transport: %s, version: %d, %s%s",
4077 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
4085 LogInfo(
"%s", new_peer_msg());
4088 if (
auto tx_relay = peer.GetTxRelay()) {
4097 tx_relay->m_tx_inventory_mutex,
4098 return tx_relay->m_tx_inventory_to_send.empty() &&
4099 tx_relay->m_next_inv_send_time == 0
s));
4109 PushPrivateBroadcastTx(pfrom);
4122 if (m_txreconciliation) {
4123 if (!peer.m_wtxid_relay || !m_txreconciliation->IsPeerRegistered(pfrom.
GetId())) {
4127 m_txreconciliation->ForgetPeer(pfrom.
GetId());
4133 const CNodeState* state =
State(pfrom.
GetId());
4135 .m_preferred = state->fPreferredDownload,
4136 .m_relay_permissions = pfrom.HasPermission(NetPermissionFlags::Relay),
4137 .m_wtxid_relay = peer.m_wtxid_relay,
4146 peer.m_prefers_headers =
true;
4151 uint8_t sendcmpct_hb{0};
4152 uint64_t sendcmpct_version{0};
4153 vRecv >> sendcmpct_hb >> sendcmpct_version;
4157 if (sendcmpct_hb > 1) {
4158 Misbehaving(peer,
"invalid sendcmpct announce field");
4166 CNodeState* nodestate =
State(pfrom.
GetId());
4167 nodestate->m_provides_cmpctblocks =
true;
4168 nodestate->m_requested_hb_cmpctblocks = sendcmpct_hb;
4185 if (!peer.m_wtxid_relay) {
4186 peer.m_wtxid_relay =
true;
4187 m_wtxid_relay_peers++;
4206 peer.m_wants_addrv2 =
true;
4223 std::string feature_id;
4227 std::vector<unsigned char> feature_data_vec;
4230 }
catch (
const std::exception&) {
4233 if (feature_id.size() < 4 || !vRecv.
empty()) {
4253 if (!m_txreconciliation) {
4254 LogDebug(
BCLog::NET,
"sendtxrcncl from peer=%d ignored, as our node does not have txreconciliation enabled\n", pfrom.
GetId());
4265 if (RejectIncomingTxs(pfrom)) {
4274 const auto* tx_relay = peer.GetTxRelay();
4275 if (!tx_relay || !
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs)) {
4281 uint32_t peer_txreconcl_version;
4282 uint64_t remote_salt;
4283 vRecv >> peer_txreconcl_version >> remote_salt;
4286 peer_txreconcl_version, remote_salt);
4318 const auto ser_params{
4326 std::vector<CAddress> vAddr;
4327 vRecv >> ser_params(vAddr);
4328 ProcessAddrs(msg_type, pfrom, peer, std::move(vAddr), interruptMsgProc);
4333 std::vector<CInv> vInv;
4337 Misbehaving(peer,
strprintf(
"inv message size = %u", vInv.size()));
4341 const bool reject_tx_invs{RejectIncomingTxs(pfrom)};
4342 std::unordered_set<uint256, SaltedUint256Hasher> seen_txids{0, m_txhash_hasher};
4343 std::unordered_set<uint256, SaltedUint256Hasher> seen_wtxids{0, m_txhash_hasher};
4347 const auto current_time{GetTime<std::chrono::microseconds>()};
4350 for (
CInv& inv : vInv) {
4351 if (interruptMsgProc)
return;
4356 if (peer.m_wtxid_relay) {
4363 const bool fAlreadyHave = AlreadyHaveBlock(inv.
hash);
4366 UpdateBlockAvailability(pfrom.
GetId(), inv.
hash);
4374 best_block = &inv.
hash;
4377 if (reject_tx_invs) {
4383 auto& seen_hashes{inv.
IsMsgWtx() ? seen_wtxids : seen_txids};
4384 if (!seen_hashes.insert(inv.
hash).second)
continue;
4386 AddKnownTx(peer, inv.
hash);
4389 const bool fAlreadyHave{m_txdownloadman.AddTxAnnouncement(pfrom.
GetId(), gtxid, current_time)};
4397 if (best_block !=
nullptr) {
4409 if (state.fSyncStarted || (!peer.m_inv_triggered_getheaders_before_sync && *best_block != m_last_block_inv_triggering_headers_sync)) {
4410 if (MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer)) {
4412 m_chainman.m_best_header->nHeight, best_block->ToString(),
4415 if (!state.fSyncStarted) {
4416 peer.m_inv_triggered_getheaders_before_sync =
true;
4420 m_last_block_inv_triggering_headers_sync = *best_block;
4429 std::vector<CInv> vInv;
4433 Misbehaving(peer,
strprintf(
"getdata message size = %u", vInv.size()));
4439 if (vInv.size() > 0) {
4444 const auto pushed_tx_opt{m_tx_for_private_broadcast.GetTxForNode(pfrom.
GetId())};
4445 if (!pushed_tx_opt) {
4456 if (vInv.size() == 1 && vInv[0].IsMsgTx() && vInv[0].hash == pushed_tx->GetHash().ToUint256()) {
4460 peer.m_ping_queued =
true;
4471 LOCK(peer.m_getdata_requests_mutex);
4472 peer.m_getdata_requests.insert(peer.m_getdata_requests.end(), vInv.begin(), vInv.end());
4473 ProcessGetData(pfrom, peer, interruptMsgProc);
4482 vRecv >> locator >> hashStop;
4498 std::shared_ptr<const CBlock> a_recent_block;
4500 LOCK(m_most_recent_block_mutex);
4501 a_recent_block = m_most_recent_block;
4504 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
4519 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4534 if (--nLimit <= 0) {
4538 WITH_LOCK(peer.m_block_inv_mutex, {peer.m_continuation_block = pindex->GetBlockHash();});
4558 for (
size_t i = 1; i < req.
indexes.size(); ++i) {
4562 std::shared_ptr<const CBlock> recent_block;
4564 LOCK(m_most_recent_block_mutex);
4565 if (m_most_recent_block_hash == req.
blockhash)
4566 recent_block = m_most_recent_block;
4570 SendBlockTransactions(pfrom, peer, *recent_block, req);
4589 if (!block_pos.IsNull()) {
4596 SendBlockTransactions(pfrom, peer, block, req);
4609 WITH_LOCK(peer.m_getdata_requests_mutex, peer.m_getdata_requests.push_back(inv));
4617 vRecv >> locator >> hashStop;
4635 if (m_chainman.
ActiveTip() ==
nullptr ||
4637 LogDebug(
BCLog::NET,
"Ignoring getheaders from peer=%d because active chain has too little work; sending empty response\n", pfrom.
GetId());
4644 CNodeState *nodestate =
State(pfrom.
GetId());
4653 if (!BlockRequestAllowed(*pindex)) {
4654 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block header that isn't in the main chain\n", __func__, pfrom.
GetId());
4667 std::vector<CBlock> vHeaders;
4668 int nLimit = m_opts.max_headers_result;
4670 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4673 if (--nLimit <= 0 || pindex->GetBlockHash() == hashStop)
4688 nodestate->pindexBestHeaderSent = pindex ? pindex : m_chainman.
ActiveChain().
Tip();
4694 if (RejectIncomingTxs(pfrom)) {
4708 const Txid& txid = ptx->GetHash();
4709 const Wtxid& wtxid = ptx->GetWitnessHash();
4712 AddKnownTx(peer, hash);
4714 if (
const auto num_broadcasted{m_tx_for_private_broadcast.Remove(ptx)}) {
4716 "network from %s; stopping private broadcast attempts",
4727 const auto& [should_validate, package_to_validate] = m_txdownloadman.ReceivedTx(pfrom.
GetId(), ptx);
4728 if (!should_validate) {
4733 if (!m_mempool.
exists(txid)) {
4734 LogInfo(
"Not relaying non-mempool transaction %s (wtxid=%s) from forcerelay peer=%d\n",
4737 LogInfo(
"Force relaying tx %s (wtxid=%s) from peer=%d\n",
4739 InitiateTxBroadcastToAll(wtxid);
4743 if (package_to_validate) {
4746 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4747 ProcessPackageResult(package_to_validate.value(), package_result);
4753 Assume(!package_to_validate.has_value());
4763 if (
auto package_to_validate{ProcessInvalidTx(pfrom.
GetId(), ptx, state,
true)}) {
4766 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4767 ProcessPackageResult(package_to_validate.value(), package_result);
4780 }
else if (m_opts.ignore_incoming_txs) {
4787 const CNodeState *nodestate =
State(pfrom.
GetId());
4788 if (!nodestate->m_provides_cmpctblocks) {
4795 vRecv >> cmpctblock;
4797 bool received_new_header =
false;
4807 MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer);
4817 received_new_header =
true;
4825 MaybePunishNodeForBlock(pfrom.
GetId(), state,
true,
"invalid header via cmpctblock");
4832 if (received_new_header) {
4833 LogBlockHeader(*pindex, pfrom,
true);
4836 bool fProcessBLOCKTXN =
false;
4840 bool fRevertToHeaderProcessing =
false;
4844 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
4845 bool fBlockReconstructed =
false;
4851 CNodeState *nodestate =
State(pfrom.
GetId());
4856 nodestate->m_last_block_announcement =
GetTime();
4862 auto range_flight = mapBlocksInFlight.equal_range(pindex->
GetBlockHash());
4863 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
4864 bool requested_block_from_this_peer{
false};
4867 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
4869 while (range_flight.first != range_flight.second) {
4870 if (range_flight.first->second.first == pfrom.
GetId()) {
4871 requested_block_from_this_peer =
true;
4874 range_flight.first++;
4884 if (requested_block_from_this_peer) {
4887 std::vector<CInv> vInv(1);
4895 if (!already_in_flight && !CanDirectFetch()) {
4903 requested_block_from_this_peer) {
4904 std::list<QueuedBlock>::iterator* queuedBlockIt =
nullptr;
4905 if (!BlockRequested(pfrom.
GetId(), *pindex, &queuedBlockIt)) {
4906 if (!(*queuedBlockIt)->partialBlock)
4919 Misbehaving(peer,
"invalid compact block");
4922 if (first_in_flight) {
4924 std::vector<CInv> vInv(1);
4935 for (
size_t i = 0; i < cmpctblock.
BlockTxCount(); i++) {
4940 fProcessBLOCKTXN =
true;
4941 }
else if (first_in_flight) {
4948 IsBlockRequestedFromOutbound(blockhash) ||
4967 ReadStatus status = tempBlock.InitData(cmpctblock, vExtraTxnForCompact);
4972 std::vector<CTransactionRef> dummy;
4974 status = tempBlock.FillBlock(*pblock, dummy,
4977 fBlockReconstructed =
true;
4981 if (requested_block_from_this_peer) {
4984 std::vector<CInv> vInv(1);
4990 fRevertToHeaderProcessing =
true;
4995 if (fProcessBLOCKTXN) {
4998 return ProcessCompactBlockTxns(pfrom, peer, txn);
5001 if (fRevertToHeaderProcessing) {
5007 return ProcessHeadersMessage(pfrom, peer, {cmpctblock.
header},
true);
5010 if (fBlockReconstructed) {
5015 mapBlockSource.emplace(pblock->GetHash(), std::make_pair(pfrom.
GetId(),
false));
5033 RemoveBlockRequest(pblock->GetHash(), std::nullopt);
5050 return ProcessCompactBlockTxns(pfrom, peer, resp);
5061 std::vector<CBlockHeader>
headers;
5065 if (nCount > m_opts.max_headers_result) {
5066 Misbehaving(peer,
strprintf(
"headers message size = %u", nCount));
5070 for (
unsigned int n = 0; n < nCount; n++) {
5075 ProcessHeadersMessage(pfrom, peer, std::move(
headers),
false);
5079 if (m_headers_presync_should_signal.exchange(
false)) {
5080 HeadersPresyncStats stats;
5082 LOCK(m_headers_presync_mutex);
5083 auto it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
5084 if (it != m_headers_presync_stats.end()) stats = it->second;
5102 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
5113 Misbehaving(peer,
"mutated block");
5118 bool forceProcessing =
false;
5119 const uint256 hash(pblock->GetHash());
5120 bool min_pow_checked =
false;
5125 forceProcessing = IsBlockRequested(hash);
5126 RemoveBlockRequest(hash, pfrom.
GetId());
5130 mapBlockSource.emplace(hash, std::make_pair(pfrom.
GetId(),
true));
5134 min_pow_checked =
true;
5137 ProcessBlock(pfrom, pblock, forceProcessing, min_pow_checked);
5154 Assume(SetupAddressRelay(pfrom, peer));
5158 if (peer.m_getaddr_recvd) {
5162 peer.m_getaddr_recvd =
true;
5164 peer.m_addrs_to_send.clear();
5165 std::vector<CAddress> vAddr;
5171 for (
const CAddress &addr : vAddr) {
5172 PushAddress(peer, addr);
5200 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5201 LOCK(tx_relay->m_tx_inventory_mutex);
5202 tx_relay->m_send_mempool =
true;
5228 ProcessPong(pfrom, peer, time_received, vRecv);
5244 Misbehaving(peer,
"too-large bloom filter");
5245 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5247 LOCK(tx_relay->m_bloom_filter_mutex);
5248 tx_relay->m_bloom_filter.reset(
new CBloomFilter(filter));
5249 tx_relay->m_relay_txs =
true;
5253 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5264 std::vector<unsigned char> vData;
5272 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5273 LOCK(tx_relay->m_bloom_filter_mutex);
5274 if (tx_relay->m_bloom_filter) {
5275 tx_relay->m_bloom_filter->insert(vData);
5281 Misbehaving(peer,
"bad filteradd message");
5292 auto tx_relay = peer.GetTxRelay();
5293 if (!tx_relay)
return;
5296 LOCK(tx_relay->m_bloom_filter_mutex);
5297 tx_relay->m_bloom_filter =
nullptr;
5298 tx_relay->m_relay_txs =
true;
5302 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5308 vRecv >> newFeeFilter;
5310 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5311 tx_relay->m_fee_filter_received = newFeeFilter;
5319 ProcessGetCFilters(pfrom, peer, vRecv);
5324 ProcessGetCFHeaders(pfrom, peer, vRecv);
5329 ProcessGetCFCheckPt(pfrom, peer, vRecv);
5334 std::vector<CInv> vInv;
5336 std::vector<GenTxid> tx_invs;
5338 for (
CInv &inv : vInv) {
5344 LOCK(m_tx_download_mutex);
5345 m_txdownloadman.ReceivedNotFound(pfrom.
GetId(), tx_invs);
5354bool PeerManagerImpl::MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer)
5357 LOCK(peer.m_misbehavior_mutex);
5360 if (!peer.m_should_discourage)
return false;
5362 peer.m_should_discourage =
false;
5367 LogWarning(
"Not punishing noban peer %d!", peer.m_id);
5373 LogWarning(
"Not punishing manually connected peer %d!", peer.m_id);
5393bool PeerManagerImpl::MaybeDisconnectForTxRelayCapacity(
CNode&
node,
const std::string& msg_type, std::optional<NodeId> protect_peer)
5395 if (!
node.IsInboundConn() || !
node.m_relays_txs)
return false;
5398 LogDebug(
BCLog::NET,
"failed to find a tx-relaying eviction candidate - connection dropped after %s message, peer=%d\n", msg_type,
node.GetId());
5399 node.fDisconnect =
true;
5403bool PeerManagerImpl::ProcessMessages(
CNode&
node, std::atomic<bool>& interruptMsgProc)
5408 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
5409 if (maybe_peer ==
nullptr)
return false;
5410 Peer& peer{*maybe_peer};
5414 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent)
return false;
5417 LOCK(peer.m_getdata_requests_mutex);
5418 if (!peer.m_getdata_requests.empty()) {
5419 ProcessGetData(
node, peer, interruptMsgProc);
5423 const bool processed_orphan = ProcessOrphanTx(peer);
5425 if (
node.fDisconnect)
5428 if (processed_orphan)
return true;
5433 LOCK(peer.m_getdata_requests_mutex);
5434 if (!peer.m_getdata_requests.empty())
return true;
5438 if (
node.fPauseSend)
return false;
5440 auto poll_result{
node.PollMessage()};
5447 bool fMoreWork = poll_result->second;
5451 node.m_addr_name.c_str(),
5452 node.ConnectionTypeAsString().c_str(),
5458 if (m_opts.capture_messages) {
5463 ProcessMessage(peer,
node,
msg.m_type,
msg.m_recv,
msg.m_time, interruptMsgProc);
5464 if (interruptMsgProc)
return false;
5466 LOCK(peer.m_getdata_requests_mutex);
5467 if (!peer.m_getdata_requests.empty()) fMoreWork =
true;
5474 LOCK(m_tx_download_mutex);
5475 if (m_txdownloadman.HaveMoreWork(peer.m_id)) fMoreWork =
true;
5476 }
catch (
const std::exception& e) {
5485void PeerManagerImpl::ConsiderEviction(
CNode& pto, Peer& peer, std::chrono::seconds time_in_seconds)
5498 if (state.pindexBestKnownBlock !=
nullptr && state.pindexBestKnownBlock->nChainWork >= m_chainman.
ActiveChain().
Tip()->
nChainWork) {
5500 if (state.m_chain_sync.m_timeout != 0
s) {
5501 state.m_chain_sync.m_timeout = 0
s;
5502 state.m_chain_sync.m_work_header =
nullptr;
5503 state.m_chain_sync.m_sent_getheaders =
false;
5505 }
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)) {
5513 state.m_chain_sync.m_work_header = m_chainman.
ActiveChain().
Tip();
5514 state.m_chain_sync.m_sent_getheaders =
false;
5515 }
else if (state.m_chain_sync.m_timeout > 0
s && time_in_seconds > state.m_chain_sync.m_timeout) {
5519 if (state.m_chain_sync.m_sent_getheaders) {
5521 LogInfo(
"Outbound peer has old chain, best known block = %s, %s", state.pindexBestKnownBlock !=
nullptr ? state.pindexBestKnownBlock->GetBlockHash().ToString() :
"<none>", pto.
DisconnectMsg());
5524 assert(state.m_chain_sync.m_work_header);
5529 MaybeSendGetHeaders(pto,
5530 GetLocator(state.m_chain_sync.m_work_header->pprev),
5532 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());
5533 state.m_chain_sync.m_sent_getheaders =
true;
5554 std::pair<NodeId, std::chrono::seconds> youngest_peer{-1, 0}, next_youngest_peer{-1, 0};
5558 if (pnode->
GetId() > youngest_peer.first) {
5559 next_youngest_peer = youngest_peer;
5560 youngest_peer.first = pnode->GetId();
5561 youngest_peer.second = pnode->m_last_block_time;
5564 NodeId to_disconnect = youngest_peer.first;
5565 if (youngest_peer.second > next_youngest_peer.second) {
5568 to_disconnect = next_youngest_peer.first;
5577 CNodeState *node_state =
State(pnode->
GetId());
5578 if (node_state ==
nullptr ||
5581 LogDebug(
BCLog::NET,
"disconnecting extra block-relay-only peer=%d (last block received at time %d)\n",
5585 LogDebug(
BCLog::NET,
"keeping block-relay-only peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5586 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), node_state->vBlocksInFlight.size());
5601 int64_t oldest_block_announcement = std::numeric_limits<int64_t>::max();
5604 AssertLockHeld(::cs_main);
5608 if (!pnode->IsFullOutboundConn() || pnode->fDisconnect) return;
5609 CNodeState *state = State(pnode->GetId());
5610 if (state == nullptr) return;
5612 if (state->m_chain_sync.m_protect) return;
5615 if (!m_connman.MultipleManualOrFullOutboundConns(pnode->addr.GetNetwork())) return;
5616 if (state->m_last_block_announcement < oldest_block_announcement || (state->m_last_block_announcement == oldest_block_announcement && pnode->GetId() > worst_peer)) {
5617 worst_peer = pnode->GetId();
5618 oldest_block_announcement = state->m_last_block_announcement;
5621 if (worst_peer != -1) {
5632 LogDebug(
BCLog::NET,
"disconnecting extra outbound peer=%d (last block announcement received at time %d)\n", pnode->
GetId(), oldest_block_announcement);
5636 LogDebug(
BCLog::NET,
"keeping outbound peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5637 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), state.vBlocksInFlight.size());
5653void PeerManagerImpl::CheckForStaleTipAndEvictPeers()
5658 auto now{GetTime<std::chrono::seconds>()};
5660 EvictExtraOutboundPeers(current_time);
5662 if (now > m_stale_tip_check_time) {
5666 LogInfo(
"Potential stale tip detected, will try using extra outbound peer (last tip update: %d seconds ago)\n",
5675 if (!m_initial_sync_finished && CanDirectFetch()) {
5677 m_initial_sync_finished =
true;
5684 peer.m_ping_nonce_sent &&
5694 bool pingSend =
false;
5696 if (peer.m_ping_queued) {
5701 if (peer.m_ping_nonce_sent == 0 && now > peer.m_ping_start.load() +
PING_INTERVAL) {
5710 }
while (
nonce == 0);
5711 peer.m_ping_queued =
false;
5712 peer.m_ping_start = now;
5714 peer.m_ping_nonce_sent =
nonce;
5718 peer.m_ping_nonce_sent = 0;
5724void PeerManagerImpl::MaybeSendAddr(
CNode&
node, Peer& peer, std::chrono::microseconds current_time)
5727 if (!peer.m_addr_relay_enabled)
return;
5729 LOCK(peer.m_addr_send_times_mutex);
5732 peer.m_next_local_addr_send < current_time) {
5739 if (peer.m_next_local_addr_send != 0us) {
5740 peer.m_addr_known->reset();
5743 CAddress local_addr{*local_service, peer.m_our_services, Now<NodeSeconds>()};
5744 if (peer.m_next_local_addr_send == 0us) {
5748 if (IsAddrCompatible(peer, local_addr)) {
5749 std::vector<CAddress> self_announcement{local_addr};
5750 if (peer.m_wants_addrv2) {
5758 PushAddress(peer, local_addr);
5765 if (current_time <= peer.m_next_addr_send)
return;
5778 bool ret = peer.m_addr_known->contains(addr.
GetKey());
5779 if (!
ret) peer.m_addr_known->insert(addr.
GetKey());
5782 peer.m_addrs_to_send.erase(std::remove_if(peer.m_addrs_to_send.begin(), peer.m_addrs_to_send.end(), addr_already_known),
5783 peer.m_addrs_to_send.end());
5786 if (peer.m_addrs_to_send.empty())
return;
5788 if (peer.m_wants_addrv2) {
5793 peer.m_addrs_to_send.clear();
5796 if (peer.m_addrs_to_send.capacity() > 40) {
5797 peer.m_addrs_to_send.shrink_to_fit();
5801void PeerManagerImpl::MaybeSendSendHeaders(
CNode&
node, Peer& peer)
5809 CNodeState &state = *
State(
node.GetId());
5810 if (state.pindexBestKnownBlock !=
nullptr &&
5817 peer.m_sent_sendheaders =
true;
5822void PeerManagerImpl::MaybeSendFeefilter(
CNode& pto, Peer& peer, std::chrono::microseconds current_time)
5824 if (m_opts.ignore_incoming_txs)
return;
5840 if (peer.m_fee_filter_sent == MAX_FILTER) {
5843 peer.m_next_send_feefilter = 0us;
5846 if (current_time > peer.m_next_send_feefilter) {
5847 CAmount filterToSend = m_fee_filter_rounder.round(currentFilter);
5850 if (filterToSend != peer.m_fee_filter_sent) {
5852 peer.m_fee_filter_sent = filterToSend;
5859 (currentFilter < 3 * peer.m_fee_filter_sent / 4 || currentFilter > 4 * peer.m_fee_filter_sent / 3)) {
5864bool PeerManagerImpl::RejectIncomingTxs(
const CNode& peer)
const
5877 const size_t nAvail{vRecv.
size()};
5878 bool bPingFinished =
false;
5879 std::string sProblem;
5881 if (nAvail >=
sizeof(
nonce)) {
5885 if (peer.m_ping_nonce_sent != 0) {
5886 if (
nonce == peer.m_ping_nonce_sent) {
5888 bPingFinished =
true;
5889 const auto ping_time = ping_end - peer.m_ping_start.load();
5890 if (ping_time.count() >= 0) {
5894 m_tx_for_private_broadcast.NodeConfirmedReception(pfrom.
GetId());
5901 sProblem =
"Timing mishap";
5905 sProblem =
"Nonce mismatch";
5908 bPingFinished =
true;
5909 sProblem =
"Nonce zero";
5913 sProblem =
"Unsolicited pong without ping";
5917 bPingFinished =
true;
5918 sProblem =
"Short payload";
5921 if (!(sProblem.empty())) {
5925 peer.m_ping_nonce_sent,
5929 if (bPingFinished) {
5930 peer.m_ping_nonce_sent = 0;
5934bool PeerManagerImpl::SetupAddressRelay(
const CNode&
node, Peer& peer)
5939 if (
node.IsBlockOnlyConn())
return false;
5944 if (
node.IsFeelerConn())
return false;
5946 if (!peer.m_addr_relay_enabled.exchange(
true)) {
5950 peer.m_addr_known = std::make_unique<CRollingBloomFilter>(5000, 0.001);
5956void PeerManagerImpl::ProcessAddrs(std::string_view msg_type,
CNode& pfrom, Peer& peer, std::vector<CAddress>&& vAddr,
const std::atomic<bool>& interruptMsgProc)
5961 if (!SetupAddressRelay(pfrom, peer)) {
5968 Misbehaving(peer,
strprintf(
"%s message size = %u", msg_type, vAddr.size()));
5973 std::vector<CAddress> vAddrOk;
5979 const auto time_diff{current_time - peer.m_addr_token_timestamp};
5983 peer.m_addr_token_timestamp = current_time;
5986 uint64_t num_proc = 0;
5987 uint64_t num_rate_limit = 0;
5988 std::shuffle(vAddr.begin(), vAddr.end(), m_rng);
5991 if (interruptMsgProc)
5995 if (peer.m_addr_token_bucket < 1.0) {
6001 peer.m_addr_token_bucket -= 1.0;
6010 addr.
nTime = std::chrono::time_point_cast<std::chrono::seconds>(current_time - 5 * 24h);
6012 AddAddressKnown(peer, addr);
6019 if (addr.
nTime > current_time - 10min && !peer.m_getaddr_sent && vAddr.size() <= 10 && addr.
IsRoutable()) {
6021 RelayAddress(pfrom.
GetId(), addr, reachable);
6025 vAddrOk.push_back(addr);
6028 peer.m_addr_processed += num_proc;
6029 peer.m_addr_rate_limited += num_rate_limit;
6030 LogDebug(
BCLog::NET,
"Received addr: %u addresses (%u processed, %u rate-limited) from peer=%d\n",
6031 vAddr.size(), num_proc, num_rate_limit, pfrom.
GetId());
6033 m_addrman.
Add(vAddrOk, pfrom.
addr, 2h);
6034 if (vAddr.size() < 1000) peer.m_getaddr_sent =
false;
6043bool PeerManagerImpl::SendMessages(
CNode&
node)
6048 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
6049 if (!maybe_peer)
return false;
6050 Peer& peer{*maybe_peer};
6055 if (MaybeDiscourageAndDisconnect(
node, peer))
return true;
6058 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent) {
6059 PushNodeVersion(
node, peer);
6060 peer.m_outbound_version_message_sent =
true;
6064 if (!
node.fSuccessfullyConnected ||
node.fDisconnect)
6068 const auto current_time{GetTime<std::chrono::microseconds>()};
6073 if (
node.IsPrivateBroadcastConn()) {
6077 node.fDisconnect =
true;
6084 node.fDisconnect =
true;
6088 MaybeSendPing(
node, peer, now);
6091 if (
node.fDisconnect)
return true;
6093 MaybeSendAddr(
node, peer, current_time);
6095 MaybeSendSendHeaders(
node, peer);
6097 ProcessInvBacklog(now);
6102 CNodeState &state = *
State(
node.GetId());
6105 if (m_chainman.m_best_header ==
nullptr) {
6112 bool sync_blocks_and_headers_from_peer =
false;
6113 if (state.fPreferredDownload) {
6114 sync_blocks_and_headers_from_peer =
true;
6115 }
else if (CanServeBlocks(peer) && !
node.IsAddrFetchConn()) {
6125 if (m_num_preferred_download_peers == 0 || mapBlocksInFlight.empty()) {
6126 sync_blocks_and_headers_from_peer =
true;
6132 if ((nSyncStarted == 0 && sync_blocks_and_headers_from_peer) || m_chainman.m_best_header->Time() >
NodeClock::now() - 24h) {
6133 const CBlockIndex* pindexStart = m_chainman.m_best_header;
6141 if (pindexStart->
pprev)
6142 pindexStart = pindexStart->
pprev;
6146 state.fSyncStarted =
true;
6170 LOCK(peer.m_block_inv_mutex);
6171 std::vector<CBlock> vHeaders;
6172 bool fRevertToInv = ((!peer.m_prefers_headers &&
6173 (!state.m_requested_hb_cmpctblocks || peer.m_blocks_for_headers_relay.size() > 1)) ||
6176 ProcessBlockAvailability(
node.GetId());
6178 if (!fRevertToInv) {
6179 bool fFoundStartingHeader =
false;
6183 for (
const uint256& hash : peer.m_blocks_for_headers_relay) {
6188 fRevertToInv =
true;
6191 if (pBestIndex !=
nullptr && pindex->
pprev != pBestIndex) {
6203 fRevertToInv =
true;
6206 pBestIndex = pindex;
6207 if (fFoundStartingHeader) {
6210 }
else if (PeerHasHeader(&state, pindex)) {
6212 }
else if (pindex->
pprev ==
nullptr || PeerHasHeader(&state, pindex->
pprev)) {
6215 fFoundStartingHeader =
true;
6220 fRevertToInv =
true;
6225 if (!fRevertToInv && !vHeaders.empty()) {
6226 if (vHeaders.size() == 1 && state.m_requested_hb_cmpctblocks) {
6230 vHeaders.front().GetHash().ToString(),
node.GetId());
6232 std::optional<CSerializedNetMsg> cached_cmpctblock_msg;
6234 LOCK(m_most_recent_block_mutex);
6235 if (m_most_recent_block_hash == pBestIndex->
GetBlockHash()) {
6239 if (cached_cmpctblock_msg.has_value()) {
6240 PushMessage(
node, std::move(cached_cmpctblock_msg.value()));
6248 state.pindexBestHeaderSent = pBestIndex;
6249 }
else if (peer.m_prefers_headers) {
6250 if (vHeaders.size() > 1) {
6253 vHeaders.front().GetHash().ToString(),
6254 vHeaders.back().GetHash().ToString(),
node.GetId());
6257 vHeaders.front().GetHash().ToString(),
node.GetId());
6260 state.pindexBestHeaderSent = pBestIndex;
6262 fRevertToInv =
true;
6268 if (!peer.m_blocks_for_headers_relay.empty()) {
6269 const uint256& hashToAnnounce = peer.m_blocks_for_headers_relay.back();
6282 if (!PeerHasHeader(&state, pindex)) {
6283 peer.m_blocks_for_inv_relay.push_back(hashToAnnounce);
6289 peer.m_blocks_for_headers_relay.clear();
6295 std::vector<CInv> vInv;
6297 LOCK(peer.m_block_inv_mutex);
6298 vInv.reserve(peer.m_blocks_for_inv_relay.size());
6301 for (
const uint256& hash : peer.m_blocks_for_inv_relay) {
6308 peer.m_blocks_for_inv_relay.clear();
6311 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
6312 LOCK(tx_relay->m_tx_inventory_mutex);
6315 if (tx_relay->m_next_inv_send_time < current_time) {
6316 fSendTrickle =
true;
6317 if (
node.IsInboundConn()) {
6326 LOCK(tx_relay->m_bloom_filter_mutex);
6327 if (!tx_relay->m_relay_txs) tx_relay->m_tx_inventory_to_send.clear();
6331 if (fSendTrickle && tx_relay->m_send_mempool) {
6332 auto vtxinfo = m_mempool.
infoAll();
6337 tx_relay->m_send_mempool =
false;
6338 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6341 tx_relay->m_tx_inventory_to_send.clear();
6343 LOCK(tx_relay->m_bloom_filter_mutex);
6345 for (
const auto& txinfo : vtxinfo) {
6346 const Txid& txid{txinfo.tx->GetHash()};
6347 const Wtxid& wtxid{txinfo.tx->GetWitnessHash()};
6348 const auto inv = peer.m_wtxid_relay ?
6353 if (txinfo.fee < filterrate.GetFee(txinfo.vsize)) {
6356 if (tx_relay->m_bloom_filter) {
6357 if (!tx_relay->m_bloom_filter->IsRelevantAndUpdate(*txinfo.tx))
continue;
6359 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6360 vInv.push_back(inv);
6372 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6375 auto& invs = tx_relay->m_tx_inventory_to_send;
6376 std::vector<CTransactionRef> res;
6378 if (invs.size() == 0)
return res;
6381 if (invs.capacity() > 2 * invs.size()) invs.shrink_to_fit();
6385 res.reserve(txiters.size());
6386 for (
auto txiter : txiters) {
6387 if (txiter->GetFee() < filterrate.GetFee(txiter->GetTxSize())) {
6390 res.push_back(txiter->GetSharedTx());
6393 tx_relay->m_last_inv_sequence = m_mempool.
GetSequence();
6397 LOCK(tx_relay->m_bloom_filter_mutex);
6398 vInv.reserve(std::min<size_t>(
MAX_INV_SZ, vInv.size() + inv_tx.size()));
6399 for (
auto& tx : inv_tx) {
6403 const auto inv = peer.m_wtxid_relay ?
6407 if (tx_relay->m_tx_inventory_known_filter.contains(inv.
hash)) {
6410 if (tx_relay->m_bloom_filter && !tx_relay->m_bloom_filter->IsRelevantAndUpdate(*tx))
continue;
6412 vInv.push_back(inv);
6417 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6425 auto stalling_timeout = m_block_stalling_timeout.load();
6426 if (state.m_stalling_since.count() && state.m_stalling_since < current_time - stalling_timeout) {
6430 LogInfo(
"Peer is stalling block download, %s",
node.DisconnectMsg());
6431 node.fDisconnect =
true;
6435 if (stalling_timeout != new_timeout && m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
6445 if (state.vBlocksInFlight.size() > 0) {
6446 QueuedBlock &queuedBlock = state.vBlocksInFlight.front();
6447 int nOtherPeersWithValidatedDownloads = m_peers_downloading_from - 1;
6449 LogInfo(
"Timeout downloading block %s, %s", queuedBlock.pindex->GetBlockHash().ToString(),
node.DisconnectMsg());
6450 node.fDisconnect =
true;
6455 if (state.fSyncStarted && peer.m_headers_sync_timeout < std::chrono::microseconds::max()) {
6457 if (m_chainman.m_best_header->Time() <=
NodeClock::now() - 24h) {
6458 if (current_time > peer.m_headers_sync_timeout && nSyncStarted == 1 && (m_num_preferred_download_peers - state.fPreferredDownload >= 1)) {
6465 LogInfo(
"Timeout downloading headers, %s",
node.DisconnectMsg());
6466 node.fDisconnect =
true;
6469 LogInfo(
"Timeout downloading headers from noban peer, not %s",
node.DisconnectMsg());
6475 state.fSyncStarted =
false;
6477 peer.m_headers_sync_timeout = 0us;
6483 peer.m_headers_sync_timeout = std::chrono::microseconds::max();
6489 ConsiderEviction(
node, peer, GetTime<std::chrono::seconds>());
6494 std::vector<CInv> vGetData;
6496 std::vector<const CBlockIndex*> vToDownload;
6498 auto get_inflight_budget = [&state]() {
6505 FindNextBlocksToDownload(peer, get_inflight_budget(), vToDownload, staller);
6506 auto historical_blocks{m_chainman.GetHistoricalBlockRange()};
6507 if (historical_blocks && !IsLimitedPeer(peer)) {
6511 TryDownloadingHistoricalBlocks(
6513 get_inflight_budget(),
6514 vToDownload, from_tip, historical_blocks->second);
6517 uint32_t nFetchFlags = GetFetchFlags(peer);
6519 BlockRequested(
node.GetId(), *pindex);
6523 if (state.vBlocksInFlight.empty() && staller != -1) {
6524 if (
State(staller)->m_stalling_since == 0us) {
6525 State(staller)->m_stalling_since = current_time;
6535 LOCK(m_tx_download_mutex);
6536 for (
const GenTxid& gtxid : m_txdownloadman.GetRequestsToSend(
node.GetId(), current_time)) {
6545 if (!vGetData.empty())
6548 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...