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, Peer& peer, 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, Peer& peer,
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, peer, *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)
3475 while (
CTransactionRef porphanTx = m_txdownloadman.GetTxToReconsider(peer.m_id)) {
3478 const Txid& orphanHash = porphanTx->GetHash();
3479 const Wtxid& orphan_wtxid = porphanTx->GetWitnessHash();
3496 ProcessInvalidTx(peer.m_id, porphanTx, state,
false);
3505bool PeerManagerImpl::PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
3507 const uint256& stop_hash, uint32_t max_height_diff,
3511 const bool supported_filter_type =
3514 if (!supported_filter_type) {
3516 static_cast<uint8_t
>(filter_type),
node.DisconnectMsg());
3517 node.fDisconnect =
true;
3526 if (!stop_index || !BlockRequestAllowed(*stop_index)) {
3529 node.fDisconnect =
true;
3534 uint32_t stop_height = stop_index->
nHeight;
3535 if (start_height > stop_height) {
3537 "start height %d and stop height %d, %s",
3538 start_height, stop_height,
node.DisconnectMsg());
3539 node.fDisconnect =
true;
3542 if (stop_height - start_height >= max_height_diff) {
3544 stop_height - start_height + 1, max_height_diff,
node.DisconnectMsg());
3545 node.fDisconnect =
true;
3550 if (!filter_index) {
3560 uint8_t filter_type_ser;
3561 uint32_t start_height;
3564 vRecv >> filter_type_ser >> start_height >> stop_hash;
3570 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3575 std::vector<BlockFilter> filters;
3577 LogDebug(
BCLog::NET,
"Failed to find block filter in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3582 for (
const auto& filter : filters) {
3589 uint8_t filter_type_ser;
3590 uint32_t start_height;
3593 vRecv >> filter_type_ser >> start_height >> stop_hash;
3599 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3605 if (start_height > 0) {
3607 stop_index->
GetAncestor(
static_cast<int>(start_height - 1));
3609 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3615 std::vector<uint256> filter_hashes;
3617 LogDebug(
BCLog::NET,
"Failed to find block filter hashes in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3631 uint8_t filter_type_ser;
3634 vRecv >> filter_type_ser >> stop_hash;
3640 if (!PrepareBlockFilterRequest(
node, peer, filter_type, 0, stop_hash,
3641 std::numeric_limits<uint32_t>::max(),
3642 stop_index, filter_index)) {
3650 for (
int i =
headers.size() - 1; i >= 0; i--) {
3655 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3669 bool new_block{
false};
3670 m_chainman.
ProcessNewBlock(block, force_processing, min_pow_checked, &new_block);
3672 node.m_last_block_time = GetTime<std::chrono::seconds>();
3677 RemoveBlockRequest(block->GetHash(), std::nullopt);
3680 mapBlockSource.erase(block->GetHash());
3684void PeerManagerImpl::ProcessCompactBlockTxns(
CNode& pfrom, Peer& peer,
const BlockTransactions& block_transactions)
3686 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
3687 bool fBlockRead{
false};
3691 auto range_flight = mapBlocksInFlight.equal_range(block_transactions.
blockhash);
3692 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
3693 bool requested_block_from_this_peer{
false};
3696 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
3698 while (range_flight.first != range_flight.second) {
3699 auto [node_id, block_it] = range_flight.first->second;
3700 if (node_id == pfrom.
GetId() && block_it->partialBlock) {
3701 requested_block_from_this_peer =
true;
3704 range_flight.first++;
3707 if (!requested_block_from_this_peer) {
3719 Misbehaving(peer,
"previous compact block reconstruction attempt failed");
3730 Misbehaving(peer,
"invalid compact block/non-matching block transactions");
3733 if (first_in_flight) {
3738 std::vector<CInv> invs;
3743 LogDebug(
BCLog::NET,
"Peer %d sent us a compact block but it failed to reconstruct, waiting on first download to complete\n", pfrom.
GetId());
3756 mapBlockSource.emplace(block_transactions.
blockhash, std::make_pair(pfrom.
GetId(),
false));
3771void PeerManagerImpl::LogBlockHeader(
const CBlockIndex& index,
const CNode& peer,
bool via_compact_block) {
3783 "Saw new %sheader hash=%s height=%d %s",
3784 via_compact_block ?
"cmpctblock " :
"",
3796void PeerManagerImpl::PushPrivateBroadcastTx(
CNode&
node)
3800 const auto opt_tx{m_tx_for_private_broadcast.PickTxForSend(
node.GetId(),
CService{node.addr})};
3803 node.fDisconnect =
true;
3809 tx->GetHash().ToString(), tx->HasWitness() ?
strprintf(
", wtxid=%s", tx->GetWitnessHash().ToString()) :
"",
3815void PeerManagerImpl::ProcessMessage(Peer& peer,
CNode& pfrom,
const std::string& msg_type,
DataStream& vRecv,
3817 const std::atomic<bool>& interruptMsgProc)
3832 uint64_t nNonce = 1;
3835 std::string cleanSubVer;
3836 int starting_height = -1;
3839 vRecv >> nVersion >> Using<CustomUintFormatter<8>>(nServices) >> nTime;
3854 LogDebug(
BCLog::NET,
"peer does not offer the expected services (%08x offered, %08x expected), %s",
3856 GetDesirableServiceFlags(nServices),
3869 if (!vRecv.
empty()) {
3877 if (!vRecv.
empty()) {
3878 std::string strSubVer;
3882 if (!vRecv.
empty()) {
3883 vRecv >> starting_height;
3903 PushNodeVersion(pfrom, peer);
3907 const int greatest_common_version = std::min(nVersion, pfrom.
AdvertisedVersion());
3912 peer.m_their_services = nServices;
3916 pfrom.cleanSubVer = cleanSubVer;
3927 (fRelay || (peer.m_our_services &
NODE_BLOOM))) {
3928 auto*
const tx_relay = peer.SetTxRelay();
3930 LOCK(tx_relay->m_bloom_filter_mutex);
3931 tx_relay->m_relay_txs = fRelay;
3937 LogDebug(
BCLog::NET,
"receive version message: %s: version %d, blocks=%d, us=%s, txrelay=%d, %s%s",
3938 cleanSubVer.empty() ?
"<no user agent>" : cleanSubVer, pfrom.
nVersion,
3940 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
3958 if (greatest_common_version >= 70016) {
3973 const auto* tx_relay = peer.GetTxRelay();
3974 if (tx_relay &&
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs) &&
3976 const uint64_t recon_salt = m_txreconciliation->PreRegisterPeer(pfrom.
GetId());
3989 if (MaybeDisconnectForTxRelayCapacity(pfrom, msg_type, pfrom.
GetId()))
return;
3997 m_num_preferred_download_peers += state->fPreferredDownload;
4003 bool send_getaddr{
false};
4005 send_getaddr = SetupAddressRelay(pfrom, peer);
4015 peer.m_getaddr_sent =
true;
4039 peer.m_time_offset =
NodeSeconds{std::chrono::seconds{nTime}} - Now<NodeSeconds>();
4043 m_outbound_time_offsets.Add(peer.m_time_offset);
4044 m_outbound_time_offsets.WarnIfOutOfSync();
4048 if (greatest_common_version <= 70012) {
4049 constexpr auto finalAlert{
"60010000000000000000000000ffffff7f00000000ffffff7ffeffff7f01ffffff7f00000000ffffff7f00ffffff7f002f555247454e543a20416c657274206b657920636f6d70726f6d697365642c2075706772616465207265717569726564004630440220653febd6410f470f6bae11cad19c48413becb1ac2c17f908fd0fd53bdc3abd5202206d0e9c96fe88d4a0f01ed9dedae2b6f9e00da94cad0fecaae66ecf689bf71b50"_hex};
4050 MakeAndPushMessage(pfrom,
"alert", finalAlert);
4073 auto new_peer_msg = [&]() {
4075 return strprintf(
"New %s peer connected: transport: %s, version: %d, %s%s",
4079 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
4087 LogInfo(
"%s", new_peer_msg());
4090 if (
auto tx_relay = peer.GetTxRelay()) {
4099 tx_relay->m_tx_inventory_mutex,
4100 return tx_relay->m_tx_inventory_to_send.empty() &&
4101 tx_relay->m_next_inv_send_time == 0
s));
4111 PushPrivateBroadcastTx(pfrom);
4124 if (m_txreconciliation) {
4125 if (!peer.m_wtxid_relay || !m_txreconciliation->IsPeerRegistered(pfrom.
GetId())) {
4129 m_txreconciliation->ForgetPeer(pfrom.
GetId());
4135 const CNodeState* state =
State(pfrom.
GetId());
4137 .m_preferred = state->fPreferredDownload,
4138 .m_relay_permissions = pfrom.HasPermission(NetPermissionFlags::Relay),
4139 .m_wtxid_relay = peer.m_wtxid_relay,
4148 peer.m_prefers_headers =
true;
4153 uint8_t sendcmpct_hb{0};
4154 uint64_t sendcmpct_version{0};
4155 vRecv >> sendcmpct_hb >> sendcmpct_version;
4159 if (sendcmpct_hb > 1) {
4160 Misbehaving(peer,
"invalid sendcmpct announce field");
4168 CNodeState* nodestate =
State(pfrom.
GetId());
4169 nodestate->m_provides_cmpctblocks =
true;
4170 nodestate->m_requested_hb_cmpctblocks = sendcmpct_hb;
4187 if (!peer.m_wtxid_relay) {
4188 peer.m_wtxid_relay =
true;
4189 m_wtxid_relay_peers++;
4208 peer.m_wants_addrv2 =
true;
4225 std::string feature_id;
4229 std::vector<unsigned char> feature_data_vec;
4232 }
catch (
const std::exception&) {
4235 if (feature_id.size() < 4 || !vRecv.
empty()) {
4255 if (!m_txreconciliation) {
4256 LogDebug(
BCLog::NET,
"sendtxrcncl from peer=%d ignored, as our node does not have txreconciliation enabled\n", pfrom.
GetId());
4267 if (RejectIncomingTxs(pfrom)) {
4276 const auto* tx_relay = peer.GetTxRelay();
4277 if (!tx_relay || !
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs)) {
4283 uint32_t peer_txreconcl_version;
4284 uint64_t remote_salt;
4285 vRecv >> peer_txreconcl_version >> remote_salt;
4288 peer_txreconcl_version, remote_salt);
4320 const auto ser_params{
4328 std::vector<CAddress> vAddr;
4329 vRecv >> ser_params(vAddr);
4330 ProcessAddrs(msg_type, pfrom, peer, std::move(vAddr), interruptMsgProc);
4335 std::vector<CInv> vInv;
4339 Misbehaving(peer,
strprintf(
"inv message size = %u", vInv.size()));
4343 const bool reject_tx_invs{RejectIncomingTxs(pfrom)};
4344 std::unordered_set<uint256, SaltedUint256Hasher> seen_txids{0, m_txhash_hasher};
4345 std::unordered_set<uint256, SaltedUint256Hasher> seen_wtxids{0, m_txhash_hasher};
4349 const auto current_time{GetTime<std::chrono::microseconds>()};
4352 for (
CInv& inv : vInv) {
4353 if (interruptMsgProc)
return;
4358 if (peer.m_wtxid_relay) {
4365 const bool fAlreadyHave = AlreadyHaveBlock(inv.
hash);
4368 UpdateBlockAvailability(pfrom.
GetId(), inv.
hash);
4376 best_block = &inv.
hash;
4379 if (reject_tx_invs) {
4385 auto& seen_hashes{inv.
IsMsgWtx() ? seen_wtxids : seen_txids};
4386 if (!seen_hashes.insert(inv.
hash).second)
continue;
4388 AddKnownTx(peer, inv.
hash);
4391 const bool fAlreadyHave{m_txdownloadman.AddTxAnnouncement(pfrom.
GetId(), gtxid, current_time)};
4399 if (best_block !=
nullptr) {
4411 if (state.fSyncStarted || (!peer.m_inv_triggered_getheaders_before_sync && *best_block != m_last_block_inv_triggering_headers_sync)) {
4412 if (MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer)) {
4414 m_chainman.m_best_header->nHeight, best_block->ToString(),
4417 if (!state.fSyncStarted) {
4418 peer.m_inv_triggered_getheaders_before_sync =
true;
4422 m_last_block_inv_triggering_headers_sync = *best_block;
4431 std::vector<CInv> vInv;
4435 Misbehaving(peer,
strprintf(
"getdata message size = %u", vInv.size()));
4441 if (vInv.size() > 0) {
4446 const auto pushed_tx_opt{m_tx_for_private_broadcast.GetTxForNode(pfrom.
GetId())};
4447 if (!pushed_tx_opt) {
4458 if (vInv.size() == 1 && vInv[0].IsMsgTx() && vInv[0].hash == pushed_tx->GetHash().ToUint256()) {
4462 peer.m_ping_queued =
true;
4473 LOCK(peer.m_getdata_requests_mutex);
4474 peer.m_getdata_requests.insert(peer.m_getdata_requests.end(), vInv.begin(), vInv.end());
4475 ProcessGetData(pfrom, peer, interruptMsgProc);
4484 vRecv >> locator >> hashStop;
4500 std::shared_ptr<const CBlock> a_recent_block;
4502 LOCK(m_most_recent_block_mutex);
4503 a_recent_block = m_most_recent_block;
4506 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
4521 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4536 if (--nLimit <= 0) {
4540 WITH_LOCK(peer.m_block_inv_mutex, {peer.m_continuation_block = pindex->GetBlockHash();});
4560 for (
size_t i = 1; i < req.
indexes.size(); ++i) {
4564 std::shared_ptr<const CBlock> recent_block;
4566 LOCK(m_most_recent_block_mutex);
4567 if (m_most_recent_block_hash == req.
blockhash)
4568 recent_block = m_most_recent_block;
4572 SendBlockTransactions(pfrom, peer, *recent_block, req);
4591 if (!block_pos.IsNull()) {
4598 SendBlockTransactions(pfrom, peer, block, req);
4611 WITH_LOCK(peer.m_getdata_requests_mutex, peer.m_getdata_requests.push_back(inv));
4619 vRecv >> locator >> hashStop;
4637 if (m_chainman.
ActiveTip() ==
nullptr ||
4639 LogDebug(
BCLog::NET,
"Ignoring getheaders from peer=%d because active chain has too little work; sending empty response\n", pfrom.
GetId());
4646 CNodeState *nodestate =
State(pfrom.
GetId());
4655 if (!BlockRequestAllowed(*pindex)) {
4656 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block header that isn't in the main chain\n", __func__, pfrom.
GetId());
4669 std::vector<CBlock> vHeaders;
4670 int nLimit = m_opts.max_headers_result;
4672 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4675 if (--nLimit <= 0 || pindex->GetBlockHash() == hashStop)
4690 nodestate->pindexBestHeaderSent = pindex ? pindex : m_chainman.
ActiveChain().
Tip();
4696 if (RejectIncomingTxs(pfrom)) {
4710 const Txid& txid = ptx->GetHash();
4711 const Wtxid& wtxid = ptx->GetWitnessHash();
4714 AddKnownTx(peer, hash);
4716 if (
const auto num_broadcasted{m_tx_for_private_broadcast.Remove(ptx)}) {
4718 "network from %s; stopping private broadcast attempts",
4729 const auto& [should_validate, package_to_validate] = m_txdownloadman.ReceivedTx(pfrom.
GetId(), ptx);
4730 if (!should_validate) {
4735 if (!m_mempool.
exists(txid)) {
4736 LogInfo(
"Not relaying non-mempool transaction %s (wtxid=%s) from forcerelay peer=%d\n",
4739 LogInfo(
"Force relaying tx %s (wtxid=%s) from peer=%d\n",
4741 InitiateTxBroadcastToAll(wtxid);
4745 if (package_to_validate) {
4748 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4749 ProcessPackageResult(package_to_validate.value(), package_result);
4755 Assume(!package_to_validate.has_value());
4765 if (
auto package_to_validate{ProcessInvalidTx(pfrom.
GetId(), ptx, state,
true)}) {
4768 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4769 ProcessPackageResult(package_to_validate.value(), package_result);
4782 }
else if (m_opts.ignore_incoming_txs) {
4789 const CNodeState *nodestate =
State(pfrom.
GetId());
4790 if (!nodestate->m_provides_cmpctblocks) {
4797 vRecv >> cmpctblock;
4799 bool received_new_header =
false;
4809 MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer);
4819 received_new_header =
true;
4827 MaybePunishNodeForBlock(pfrom.
GetId(), state,
true,
"invalid header via cmpctblock");
4834 if (received_new_header) {
4835 LogBlockHeader(*pindex, pfrom,
true);
4838 bool fProcessBLOCKTXN =
false;
4842 bool fRevertToHeaderProcessing =
false;
4846 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
4847 bool fBlockReconstructed =
false;
4853 CNodeState *nodestate =
State(pfrom.
GetId());
4858 nodestate->m_last_block_announcement =
GetTime();
4864 auto range_flight = mapBlocksInFlight.equal_range(pindex->
GetBlockHash());
4865 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
4866 bool requested_block_from_this_peer{
false};
4869 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
4871 while (range_flight.first != range_flight.second) {
4872 if (range_flight.first->second.first == pfrom.
GetId()) {
4873 requested_block_from_this_peer =
true;
4876 range_flight.first++;
4886 if (requested_block_from_this_peer) {
4889 std::vector<CInv> vInv(1);
4897 if (!already_in_flight && !CanDirectFetch()) {
4905 requested_block_from_this_peer) {
4906 std::list<QueuedBlock>::iterator* queuedBlockIt =
nullptr;
4907 if (!BlockRequested(pfrom.
GetId(), *pindex, &queuedBlockIt)) {
4908 if (!(*queuedBlockIt)->partialBlock)
4921 Misbehaving(peer,
"invalid compact block");
4924 if (first_in_flight) {
4926 std::vector<CInv> vInv(1);
4937 for (
size_t i = 0; i < cmpctblock.
BlockTxCount(); i++) {
4942 fProcessBLOCKTXN =
true;
4943 }
else if (first_in_flight) {
4950 IsBlockRequestedFromOutbound(blockhash) ||
4969 ReadStatus status = tempBlock.InitData(cmpctblock, vExtraTxnForCompact);
4974 std::vector<CTransactionRef> dummy;
4976 status = tempBlock.FillBlock(*pblock, dummy,
4979 fBlockReconstructed =
true;
4983 if (requested_block_from_this_peer) {
4986 std::vector<CInv> vInv(1);
4992 fRevertToHeaderProcessing =
true;
4997 if (fProcessBLOCKTXN) {
5000 return ProcessCompactBlockTxns(pfrom, peer, txn);
5003 if (fRevertToHeaderProcessing) {
5009 return ProcessHeadersMessage(pfrom, peer, {cmpctblock.
header},
true);
5012 if (fBlockReconstructed) {
5017 mapBlockSource.emplace(pblock->GetHash(), std::make_pair(pfrom.
GetId(),
false));
5035 RemoveBlockRequest(pblock->GetHash(), std::nullopt);
5052 return ProcessCompactBlockTxns(pfrom, peer, resp);
5063 std::vector<CBlockHeader>
headers;
5067 if (nCount > m_opts.max_headers_result) {
5068 Misbehaving(peer,
strprintf(
"headers message size = %u", nCount));
5072 for (
unsigned int n = 0; n < nCount; n++) {
5077 ProcessHeadersMessage(pfrom, peer, std::move(
headers),
false);
5081 if (m_headers_presync_should_signal.exchange(
false)) {
5082 HeadersPresyncStats stats;
5084 LOCK(m_headers_presync_mutex);
5085 auto it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
5086 if (it != m_headers_presync_stats.end()) stats = it->second;
5104 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
5115 Misbehaving(peer,
"mutated block");
5120 bool forceProcessing =
false;
5121 const uint256 hash(pblock->GetHash());
5122 bool min_pow_checked =
false;
5127 forceProcessing = IsBlockRequested(hash);
5128 RemoveBlockRequest(hash, pfrom.
GetId());
5132 mapBlockSource.emplace(hash, std::make_pair(pfrom.
GetId(),
true));
5136 min_pow_checked =
true;
5139 ProcessBlock(pfrom, pblock, forceProcessing, min_pow_checked);
5156 Assume(SetupAddressRelay(pfrom, peer));
5160 if (peer.m_getaddr_recvd) {
5164 peer.m_getaddr_recvd =
true;
5166 peer.m_addrs_to_send.clear();
5167 std::vector<CAddress> vAddr;
5173 for (
const CAddress &addr : vAddr) {
5174 PushAddress(peer, addr);
5202 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5203 LOCK(tx_relay->m_tx_inventory_mutex);
5204 tx_relay->m_send_mempool =
true;
5230 ProcessPong(pfrom, peer, time_received, vRecv);
5246 Misbehaving(peer,
"too-large bloom filter");
5247 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5249 LOCK(tx_relay->m_bloom_filter_mutex);
5250 tx_relay->m_bloom_filter.reset(
new CBloomFilter(filter));
5251 tx_relay->m_relay_txs =
true;
5255 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5266 std::vector<unsigned char> vData;
5274 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5275 LOCK(tx_relay->m_bloom_filter_mutex);
5276 if (tx_relay->m_bloom_filter) {
5277 tx_relay->m_bloom_filter->insert(vData);
5283 Misbehaving(peer,
"bad filteradd message");
5294 auto tx_relay = peer.GetTxRelay();
5295 if (!tx_relay)
return;
5298 LOCK(tx_relay->m_bloom_filter_mutex);
5299 tx_relay->m_bloom_filter =
nullptr;
5300 tx_relay->m_relay_txs =
true;
5304 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5310 vRecv >> newFeeFilter;
5312 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5313 tx_relay->m_fee_filter_received = newFeeFilter;
5321 ProcessGetCFilters(pfrom, peer, vRecv);
5326 ProcessGetCFHeaders(pfrom, peer, vRecv);
5331 ProcessGetCFCheckPt(pfrom, peer, vRecv);
5336 std::vector<CInv> vInv;
5338 std::vector<GenTxid> tx_invs;
5340 for (
CInv &inv : vInv) {
5346 LOCK(m_tx_download_mutex);
5347 m_txdownloadman.ReceivedNotFound(pfrom.
GetId(), tx_invs);
5356bool PeerManagerImpl::MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer)
5359 LOCK(peer.m_misbehavior_mutex);
5362 if (!peer.m_should_discourage)
return false;
5364 peer.m_should_discourage =
false;
5369 LogWarning(
"Not punishing noban peer %d!", peer.m_id);
5375 LogWarning(
"Not punishing manually connected peer %d!", peer.m_id);
5395bool PeerManagerImpl::MaybeDisconnectForTxRelayCapacity(
CNode&
node,
const std::string& msg_type, std::optional<NodeId> protect_peer)
5397 if (!
node.IsInboundConn() || !
node.m_relays_txs)
return false;
5400 LogDebug(
BCLog::NET,
"failed to find a tx-relaying eviction candidate - connection dropped after %s message, peer=%d\n", msg_type,
node.GetId());
5401 node.fDisconnect =
true;
5405bool PeerManagerImpl::ProcessMessages(
CNode&
node, std::atomic<bool>& interruptMsgProc)
5410 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
5411 if (maybe_peer ==
nullptr)
return false;
5412 Peer& peer{*maybe_peer};
5416 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent)
return false;
5419 LOCK(peer.m_getdata_requests_mutex);
5420 if (!peer.m_getdata_requests.empty()) {
5421 ProcessGetData(
node, peer, interruptMsgProc);
5425 const bool processed_orphan = ProcessOrphanTx(peer);
5427 if (
node.fDisconnect)
5430 if (processed_orphan)
return true;
5435 LOCK(peer.m_getdata_requests_mutex);
5436 if (!peer.m_getdata_requests.empty())
return true;
5440 if (
node.fPauseSend)
return false;
5442 auto poll_result{
node.PollMessage()};
5449 bool fMoreWork = poll_result->second;
5453 node.m_addr_name.c_str(),
5454 node.ConnectionTypeAsString().c_str(),
5460 if (m_opts.capture_messages) {
5465 ProcessMessage(peer,
node,
msg.m_type,
msg.m_recv,
msg.m_time, interruptMsgProc);
5466 if (interruptMsgProc)
return false;
5468 LOCK(peer.m_getdata_requests_mutex);
5469 if (!peer.m_getdata_requests.empty()) fMoreWork =
true;
5476 LOCK(m_tx_download_mutex);
5477 if (m_txdownloadman.HaveMoreWork(peer.m_id)) fMoreWork =
true;
5478 }
catch (
const std::exception& e) {
5487void PeerManagerImpl::ConsiderEviction(
CNode& pto, Peer& peer, std::chrono::seconds time_in_seconds)
5500 if (state.pindexBestKnownBlock !=
nullptr && state.pindexBestKnownBlock->nChainWork >= m_chainman.
ActiveChain().
Tip()->
nChainWork) {
5502 if (state.m_chain_sync.m_timeout != 0
s) {
5503 state.m_chain_sync.m_timeout = 0
s;
5504 state.m_chain_sync.m_work_header =
nullptr;
5505 state.m_chain_sync.m_sent_getheaders =
false;
5507 }
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)) {
5515 state.m_chain_sync.m_work_header = m_chainman.
ActiveChain().
Tip();
5516 state.m_chain_sync.m_sent_getheaders =
false;
5517 }
else if (state.m_chain_sync.m_timeout > 0
s && time_in_seconds > state.m_chain_sync.m_timeout) {
5521 if (state.m_chain_sync.m_sent_getheaders) {
5523 LogInfo(
"Outbound peer has old chain, best known block = %s, %s", state.pindexBestKnownBlock !=
nullptr ? state.pindexBestKnownBlock->GetBlockHash().ToString() :
"<none>", pto.
DisconnectMsg());
5526 assert(state.m_chain_sync.m_work_header);
5531 MaybeSendGetHeaders(pto,
5532 GetLocator(state.m_chain_sync.m_work_header->pprev),
5534 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());
5535 state.m_chain_sync.m_sent_getheaders =
true;
5556 std::pair<NodeId, std::chrono::seconds> youngest_peer{-1, 0}, next_youngest_peer{-1, 0};
5560 if (pnode->
GetId() > youngest_peer.first) {
5561 next_youngest_peer = youngest_peer;
5562 youngest_peer.first = pnode->GetId();
5563 youngest_peer.second = pnode->m_last_block_time;
5566 NodeId to_disconnect = youngest_peer.first;
5567 if (youngest_peer.second > next_youngest_peer.second) {
5570 to_disconnect = next_youngest_peer.first;
5579 CNodeState *node_state =
State(pnode->
GetId());
5580 if (node_state ==
nullptr ||
5583 LogDebug(
BCLog::NET,
"disconnecting extra block-relay-only peer=%d (last block received at time %d)\n",
5587 LogDebug(
BCLog::NET,
"keeping block-relay-only peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5588 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), node_state->vBlocksInFlight.size());
5603 int64_t oldest_block_announcement = std::numeric_limits<int64_t>::max();
5606 AssertLockHeld(::cs_main);
5610 if (!pnode->IsFullOutboundConn() || pnode->fDisconnect) return;
5611 CNodeState *state = State(pnode->GetId());
5612 if (state == nullptr) return;
5614 if (state->m_chain_sync.m_protect) return;
5617 if (!m_connman.MultipleManualOrFullOutboundConns(pnode->addr.GetNetwork())) return;
5618 if (state->m_last_block_announcement < oldest_block_announcement || (state->m_last_block_announcement == oldest_block_announcement && pnode->GetId() > worst_peer)) {
5619 worst_peer = pnode->GetId();
5620 oldest_block_announcement = state->m_last_block_announcement;
5623 if (worst_peer != -1) {
5634 LogDebug(
BCLog::NET,
"disconnecting extra outbound peer=%d (last block announcement received at time %d)\n", pnode->
GetId(), oldest_block_announcement);
5638 LogDebug(
BCLog::NET,
"keeping outbound peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5639 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), state.vBlocksInFlight.size());
5655void PeerManagerImpl::CheckForStaleTipAndEvictPeers()
5660 auto now{GetTime<std::chrono::seconds>()};
5662 EvictExtraOutboundPeers(current_time);
5664 if (now > m_stale_tip_check_time) {
5668 LogInfo(
"Potential stale tip detected, will try using extra outbound peer (last tip update: %d seconds ago)\n",
5677 if (!m_initial_sync_finished && CanDirectFetch()) {
5679 m_initial_sync_finished =
true;
5686 peer.m_ping_nonce_sent &&
5696 bool pingSend =
false;
5698 if (peer.m_ping_queued) {
5703 if (peer.m_ping_nonce_sent == 0 && now > peer.m_ping_start.load() +
PING_INTERVAL) {
5712 }
while (
nonce == 0);
5713 peer.m_ping_queued =
false;
5714 peer.m_ping_start = now;
5716 peer.m_ping_nonce_sent =
nonce;
5720 peer.m_ping_nonce_sent = 0;
5726void PeerManagerImpl::MaybeSendAddr(
CNode&
node, Peer& peer, std::chrono::microseconds current_time)
5729 if (!peer.m_addr_relay_enabled)
return;
5731 LOCK(peer.m_addr_send_times_mutex);
5734 peer.m_next_local_addr_send < current_time) {
5741 if (peer.m_next_local_addr_send != 0us) {
5742 peer.m_addr_known->reset();
5745 CAddress local_addr{*local_service, peer.m_our_services, Now<NodeSeconds>()};
5746 if (peer.m_next_local_addr_send == 0us) {
5750 if (IsAddrCompatible(peer, local_addr)) {
5751 std::vector<CAddress> self_announcement{local_addr};
5752 if (peer.m_wants_addrv2) {
5760 PushAddress(peer, local_addr);
5767 if (current_time <= peer.m_next_addr_send)
return;
5780 bool ret = peer.m_addr_known->contains(addr.
GetKey());
5781 if (!
ret) peer.m_addr_known->insert(addr.
GetKey());
5784 peer.m_addrs_to_send.erase(std::remove_if(peer.m_addrs_to_send.begin(), peer.m_addrs_to_send.end(), addr_already_known),
5785 peer.m_addrs_to_send.end());
5788 if (peer.m_addrs_to_send.empty())
return;
5790 if (peer.m_wants_addrv2) {
5795 peer.m_addrs_to_send.clear();
5798 if (peer.m_addrs_to_send.capacity() > 40) {
5799 peer.m_addrs_to_send.shrink_to_fit();
5803void PeerManagerImpl::MaybeSendSendHeaders(
CNode&
node, Peer& peer)
5811 CNodeState &state = *
State(
node.GetId());
5812 if (state.pindexBestKnownBlock !=
nullptr &&
5819 peer.m_sent_sendheaders =
true;
5824void PeerManagerImpl::MaybeSendFeefilter(
CNode& pto, Peer& peer, std::chrono::microseconds current_time)
5826 if (m_opts.ignore_incoming_txs)
return;
5842 if (peer.m_fee_filter_sent == MAX_FILTER) {
5845 peer.m_next_send_feefilter = 0us;
5848 if (current_time > peer.m_next_send_feefilter) {
5849 CAmount filterToSend = m_fee_filter_rounder.round(currentFilter);
5852 if (filterToSend != peer.m_fee_filter_sent) {
5854 peer.m_fee_filter_sent = filterToSend;
5861 (currentFilter < 3 * peer.m_fee_filter_sent / 4 || currentFilter > 4 * peer.m_fee_filter_sent / 3)) {
5866bool PeerManagerImpl::RejectIncomingTxs(
const CNode& peer)
const
5879 const size_t nAvail{vRecv.
size()};
5880 bool bPingFinished =
false;
5881 std::string sProblem;
5883 if (nAvail >=
sizeof(
nonce)) {
5887 if (peer.m_ping_nonce_sent != 0) {
5888 if (
nonce == peer.m_ping_nonce_sent) {
5890 bPingFinished =
true;
5891 const auto ping_time = ping_end - peer.m_ping_start.load();
5892 if (ping_time.count() >= 0) {
5896 m_tx_for_private_broadcast.NodeConfirmedReception(pfrom.
GetId());
5903 sProblem =
"Timing mishap";
5907 sProblem =
"Nonce mismatch";
5910 bPingFinished =
true;
5911 sProblem =
"Nonce zero";
5915 sProblem =
"Unsolicited pong without ping";
5919 bPingFinished =
true;
5920 sProblem =
"Short payload";
5923 if (!(sProblem.empty())) {
5927 peer.m_ping_nonce_sent,
5931 if (bPingFinished) {
5932 peer.m_ping_nonce_sent = 0;
5936bool PeerManagerImpl::SetupAddressRelay(
const CNode&
node, Peer& peer)
5941 if (
node.IsBlockOnlyConn())
return false;
5946 if (
node.IsFeelerConn())
return false;
5948 if (!peer.m_addr_relay_enabled.exchange(
true)) {
5952 peer.m_addr_known = std::make_unique<CRollingBloomFilter>(5000, 0.001);
5958void PeerManagerImpl::ProcessAddrs(std::string_view msg_type,
CNode& pfrom, Peer& peer, std::vector<CAddress>&& vAddr,
const std::atomic<bool>& interruptMsgProc)
5963 if (!SetupAddressRelay(pfrom, peer)) {
5970 Misbehaving(peer,
strprintf(
"%s message size = %u", msg_type, vAddr.size()));
5975 std::vector<CAddress> vAddrOk;
5981 const auto time_diff{current_time - peer.m_addr_token_timestamp};
5985 peer.m_addr_token_timestamp = current_time;
5988 uint64_t num_proc = 0;
5989 uint64_t num_rate_limit = 0;
5990 std::shuffle(vAddr.begin(), vAddr.end(), m_rng);
5993 if (interruptMsgProc)
5997 if (peer.m_addr_token_bucket < 1.0) {
6003 peer.m_addr_token_bucket -= 1.0;
6012 addr.
nTime = std::chrono::time_point_cast<std::chrono::seconds>(current_time - 5 * 24h);
6014 AddAddressKnown(peer, addr);
6021 if (addr.
nTime > current_time - 10min && !peer.m_getaddr_sent && vAddr.size() <= 10 && addr.
IsRoutable()) {
6023 RelayAddress(pfrom.
GetId(), addr, reachable);
6027 vAddrOk.push_back(addr);
6030 peer.m_addr_processed += num_proc;
6031 peer.m_addr_rate_limited += num_rate_limit;
6032 LogDebug(
BCLog::NET,
"Received addr: %u addresses (%u processed, %u rate-limited) from peer=%d\n",
6033 vAddr.size(), num_proc, num_rate_limit, pfrom.
GetId());
6035 m_addrman.
Add(vAddrOk, pfrom.
addr, 2h);
6036 if (vAddr.size() < 1000) peer.m_getaddr_sent =
false;
6045bool PeerManagerImpl::SendMessages(
CNode&
node)
6050 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
6051 if (!maybe_peer)
return false;
6052 Peer& peer{*maybe_peer};
6057 if (MaybeDiscourageAndDisconnect(
node, peer))
return true;
6060 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent) {
6061 PushNodeVersion(
node, peer);
6062 peer.m_outbound_version_message_sent =
true;
6066 if (!
node.fSuccessfullyConnected ||
node.fDisconnect)
6070 const auto current_time{GetTime<std::chrono::microseconds>()};
6075 if (
node.IsPrivateBroadcastConn()) {
6079 node.fDisconnect =
true;
6086 node.fDisconnect =
true;
6090 MaybeSendPing(
node, peer, now);
6093 if (
node.fDisconnect)
return true;
6095 MaybeSendAddr(
node, peer, current_time);
6097 MaybeSendSendHeaders(
node, peer);
6099 ProcessInvBacklog(now);
6104 CNodeState &state = *
State(
node.GetId());
6107 if (m_chainman.m_best_header ==
nullptr) {
6114 bool sync_blocks_and_headers_from_peer =
false;
6115 if (state.fPreferredDownload) {
6116 sync_blocks_and_headers_from_peer =
true;
6117 }
else if (CanServeBlocks(peer) && !
node.IsAddrFetchConn()) {
6127 if (m_num_preferred_download_peers == 0 || mapBlocksInFlight.empty()) {
6128 sync_blocks_and_headers_from_peer =
true;
6134 if ((nSyncStarted == 0 && sync_blocks_and_headers_from_peer) || m_chainman.m_best_header->Time() >
NodeClock::now() - 24h) {
6135 const CBlockIndex* pindexStart = m_chainman.m_best_header;
6143 if (pindexStart->
pprev)
6144 pindexStart = pindexStart->
pprev;
6148 state.fSyncStarted =
true;
6172 LOCK(peer.m_block_inv_mutex);
6173 std::vector<CBlock> vHeaders;
6174 bool fRevertToInv = ((!peer.m_prefers_headers &&
6175 (!state.m_requested_hb_cmpctblocks || peer.m_blocks_for_headers_relay.size() > 1)) ||
6178 ProcessBlockAvailability(
node.GetId());
6180 if (!fRevertToInv) {
6181 bool fFoundStartingHeader =
false;
6185 for (
const uint256& hash : peer.m_blocks_for_headers_relay) {
6190 fRevertToInv =
true;
6193 if (pBestIndex !=
nullptr && pindex->
pprev != pBestIndex) {
6205 fRevertToInv =
true;
6208 pBestIndex = pindex;
6209 if (fFoundStartingHeader) {
6212 }
else if (PeerHasHeader(&state, pindex)) {
6214 }
else if (pindex->
pprev ==
nullptr || PeerHasHeader(&state, pindex->
pprev)) {
6217 fFoundStartingHeader =
true;
6222 fRevertToInv =
true;
6227 if (!fRevertToInv && !vHeaders.empty()) {
6228 if (vHeaders.size() == 1 && state.m_requested_hb_cmpctblocks) {
6232 vHeaders.front().GetHash().ToString(),
node.GetId());
6234 std::optional<CSerializedNetMsg> cached_cmpctblock_msg;
6236 LOCK(m_most_recent_block_mutex);
6237 if (m_most_recent_block_hash == pBestIndex->
GetBlockHash()) {
6241 if (cached_cmpctblock_msg.has_value()) {
6242 PushMessage(
node, std::move(cached_cmpctblock_msg.value()));
6250 state.pindexBestHeaderSent = pBestIndex;
6251 }
else if (peer.m_prefers_headers) {
6252 if (vHeaders.size() > 1) {
6255 vHeaders.front().GetHash().ToString(),
6256 vHeaders.back().GetHash().ToString(),
node.GetId());
6259 vHeaders.front().GetHash().ToString(),
node.GetId());
6262 state.pindexBestHeaderSent = pBestIndex;
6264 fRevertToInv =
true;
6270 if (!peer.m_blocks_for_headers_relay.empty()) {
6271 const uint256& hashToAnnounce = peer.m_blocks_for_headers_relay.back();
6284 if (!PeerHasHeader(&state, pindex)) {
6285 peer.m_blocks_for_inv_relay.push_back(hashToAnnounce);
6291 peer.m_blocks_for_headers_relay.clear();
6297 std::vector<CInv> vInv;
6299 LOCK(peer.m_block_inv_mutex);
6300 vInv.reserve(peer.m_blocks_for_inv_relay.size());
6303 for (
const uint256& hash : peer.m_blocks_for_inv_relay) {
6310 peer.m_blocks_for_inv_relay.clear();
6313 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
6314 LOCK(tx_relay->m_tx_inventory_mutex);
6317 if (tx_relay->m_next_inv_send_time < current_time) {
6318 fSendTrickle =
true;
6319 if (
node.IsInboundConn()) {
6328 LOCK(tx_relay->m_bloom_filter_mutex);
6329 if (!tx_relay->m_relay_txs) tx_relay->m_tx_inventory_to_send.clear();
6333 if (fSendTrickle && tx_relay->m_send_mempool) {
6334 auto vtxinfo = m_mempool.
infoAll();
6339 tx_relay->m_send_mempool =
false;
6340 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6343 tx_relay->m_tx_inventory_to_send.clear();
6345 LOCK(tx_relay->m_bloom_filter_mutex);
6347 for (
const auto& txinfo : vtxinfo) {
6348 const Txid& txid{txinfo.tx->GetHash()};
6349 const Wtxid& wtxid{txinfo.tx->GetWitnessHash()};
6350 const auto inv = peer.m_wtxid_relay ?
6355 if (txinfo.fee < filterrate.GetFee(txinfo.vsize)) {
6358 if (tx_relay->m_bloom_filter) {
6359 if (!tx_relay->m_bloom_filter->IsRelevantAndUpdate(*txinfo.tx))
continue;
6361 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6362 vInv.push_back(inv);
6374 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6377 auto& invs = tx_relay->m_tx_inventory_to_send;
6378 std::vector<CTransactionRef> res;
6380 if (invs.size() == 0)
return res;
6383 if (invs.capacity() > 2 * invs.size()) invs.shrink_to_fit();
6387 res.reserve(txiters.size());
6388 for (
auto txiter : txiters) {
6389 if (txiter->GetFee() < filterrate.GetFee(txiter->GetTxSize())) {
6392 res.push_back(txiter->GetSharedTx());
6395 tx_relay->m_last_inv_sequence = m_mempool.
GetSequence();
6399 LOCK(tx_relay->m_bloom_filter_mutex);
6400 vInv.reserve(std::min<size_t>(
MAX_INV_SZ, vInv.size() + inv_tx.size()));
6401 for (
auto& tx : inv_tx) {
6405 const auto inv = peer.m_wtxid_relay ?
6409 if (tx_relay->m_tx_inventory_known_filter.contains(inv.
hash)) {
6412 if (tx_relay->m_bloom_filter && !tx_relay->m_bloom_filter->IsRelevantAndUpdate(*tx))
continue;
6414 vInv.push_back(inv);
6419 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6427 auto stalling_timeout = m_block_stalling_timeout.load();
6428 if (state.m_stalling_since.count() && state.m_stalling_since < current_time - stalling_timeout) {
6432 LogInfo(
"Peer is stalling block download, %s",
node.DisconnectMsg());
6433 node.fDisconnect =
true;
6437 if (stalling_timeout != new_timeout && m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
6447 if (state.vBlocksInFlight.size() > 0) {
6448 QueuedBlock &queuedBlock = state.vBlocksInFlight.front();
6449 int nOtherPeersWithValidatedDownloads = m_peers_downloading_from - 1;
6451 LogInfo(
"Timeout downloading block %s, %s", queuedBlock.pindex->GetBlockHash().ToString(),
node.DisconnectMsg());
6452 node.fDisconnect =
true;
6457 if (state.fSyncStarted && peer.m_headers_sync_timeout < std::chrono::microseconds::max()) {
6459 if (m_chainman.m_best_header->Time() <=
NodeClock::now() - 24h) {
6460 if (current_time > peer.m_headers_sync_timeout && nSyncStarted == 1 && (m_num_preferred_download_peers - state.fPreferredDownload >= 1)) {
6467 LogInfo(
"Timeout downloading headers, %s",
node.DisconnectMsg());
6468 node.fDisconnect =
true;
6471 LogInfo(
"Timeout downloading headers from noban peer, not %s",
node.DisconnectMsg());
6477 state.fSyncStarted =
false;
6479 peer.m_headers_sync_timeout = 0us;
6485 peer.m_headers_sync_timeout = std::chrono::microseconds::max();
6491 ConsiderEviction(
node, peer, GetTime<std::chrono::seconds>());
6496 std::vector<CInv> vGetData;
6498 std::vector<const CBlockIndex*> vToDownload;
6500 auto get_inflight_budget = [&state]() {
6507 FindNextBlocksToDownload(peer, get_inflight_budget(), vToDownload, staller);
6508 auto historical_blocks{m_chainman.GetHistoricalBlockRange()};
6509 if (historical_blocks && !IsLimitedPeer(peer)) {
6513 TryDownloadingHistoricalBlocks(
6515 get_inflight_budget(),
6516 vToDownload, from_tip, historical_blocks->second);
6519 uint32_t nFetchFlags = GetFetchFlags(peer);
6521 BlockRequested(
node.GetId(), *pindex);
6525 if (state.vBlocksInFlight.empty() && staller != -1) {
6526 if (
State(staller)->m_stalling_since == 0us) {
6527 State(staller)->m_stalling_since = current_time;
6537 LOCK(m_tx_download_mutex);
6538 for (
const GenTxid& gtxid : m_txdownloadman.GetRequestsToSend(
node.GetId(), current_time)) {
6547 if (!vGetData.empty())
6550 MaybeSendFeefilter(
node, peer, current_time);
static 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.
static 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; no change.
@ Added
The transaction was newly added.
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...
static 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
static const unsigned int MAX_SUBVERSION_LENGTH
Maximum length of the user agent string in version message.
static 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.
static constexpr uint64_t CMPCTBLOCKS_VERSION
The compactblocks version we support.
static const unsigned int MAX_CMPCTBLOCKS_INFLIGHT_PER_BLOCK
Maximum number of outstanding CMPCTBLOCK requests for the same block.
ReachableNets g_reachable_nets
bool IsProxy(const CNetAddr &addr)
static constexpr unsigned int DEFAULT_MIN_RELAY_TX_FEE
Default for -minrelaytxfee, minimum relay fee for transactions.
static constexpr TransactionSerParams TX_NO_WITNESS
static 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.
const uint32_t MSG_WITNESS_FLAG
getdata message type flags
static constexpr size_t MAX_FEATUREDATA_LENGTH
@ MSG_WTX
Defined in BIP 339.
@ MSG_CMPCT_BLOCK
Defined in BIP152.
@ MSG_WITNESS_BLOCK
Defined in BIP144.
ServiceFlags
nServices flags
static 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.
static const int WTXID_RELAY_VERSION
"wtxidrelay" message type for wtxid-based relay starts with this version
static const int FEATURE_VERSION
"feature" message type for feature negotiation starts with this version
static const int SHORT_IDS_BLOCKS_VERSION
short-id-based block download starts with this version
static const int SENDHEADERS_VERSION
"sendheaders" message type and announcing blocks with headers starts with this version
static const int FEEFILTER_VERSION
"feefilter" tells peers to filter invs to you by fee starts with this version
static const int MIN_PEER_PROTO_VERSION
disconnect from peers older than this proto version
static const int INVALID_CB_NO_BAN_VERSION
not banning for invalid compact blocks starts with this version
static const int BIP0031_VERSION
BIP 0031, pong message, is enabled for all versions AFTER this one.
static const 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
static 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.
static const unsigned int MIN_BLOCKS_TO_KEEP
Block files containing a block-height within MIN_BLOCKS_TO_KEEP of ActiveChain().Tip() will not be pr...
@ UNVALIDATED
Blocks after an assumeutxo snapshot have been validated but the snapshot itself has not been validate...