76#include <initializer_list>
89#include <unordered_set>
216 std::unique_ptr<PartiallyDownloadedBlock> partialBlock;
250 std::atomic<ServiceFlags> m_their_services{
NODE_NONE};
253 const bool m_is_inbound;
256 Mutex m_misbehavior_mutex;
258 bool m_should_discourage
GUARDED_BY(m_misbehavior_mutex){
false};
261 Mutex m_block_inv_mutex;
265 std::vector<uint256> m_blocks_for_inv_relay
GUARDED_BY(m_block_inv_mutex);
269 std::vector<uint256> m_blocks_for_headers_relay
GUARDED_BY(m_block_inv_mutex);
280 std::atomic<uint64_t> m_ping_nonce_sent{0};
284 std::atomic<bool> m_ping_queued{
false};
287 std::atomic<bool> m_wtxid_relay{
false};
299 bool m_relay_txs
GUARDED_BY(m_bloom_filter_mutex){
false};
301 std::unique_ptr<CBloomFilter> m_bloom_filter
PT_GUARDED_BY(m_bloom_filter_mutex)
GUARDED_BY(m_bloom_filter_mutex){
nullptr};
312 std::vector<Wtxid> m_tx_inventory_to_send
GUARDED_BY(m_tx_inventory_mutex);
316 bool m_send_mempool
GUARDED_BY(m_tx_inventory_mutex){
false};
319 std::chrono::microseconds m_next_inv_send_time
GUARDED_BY(m_tx_inventory_mutex){0};
322 uint64_t m_last_inv_sequence
GUARDED_BY(m_tx_inventory_mutex){1};
325 std::atomic<CAmount> m_fee_filter_received{0};
331 LOCK(m_tx_relay_mutex);
333 m_tx_relay = std::make_unique<Peer::TxRelay>();
334 return m_tx_relay.get();
339 return WITH_LOCK(m_tx_relay_mutex,
return m_tx_relay.get());
368 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 std::chrono::microseconds m_block_download_paused_until{0us};
460 bool fPreferredDownload{
false};
462 bool m_requested_hb_cmpctblocks{
false};
464 bool m_provides_cmpctblocks{
false};
490 struct ChainSyncTimeoutState {
492 std::chrono::seconds m_timeout{0
s};
496 bool m_sent_getheaders{
false};
498 bool m_protect{
false};
501 ChainSyncTimeoutState m_chain_sync;
507struct InvToSendBucket {
508 const double count_floor{0};
509 std::vector<Wtxid> backlog;
526 static constexpr double SIZE_INIT{12'000'000};
527 static constexpr double SIZE_CAP{50'000'000};
528 static constexpr double SIZE_REFILL{20'000};
530 static constexpr double INBOUND_COUNT_SECONDS{30};
532 InvToSendBucket(
unsigned int rate,
double mult)
534 size_bucket(SIZE_REFILL * mult, SIZE_INIT, SIZE_CAP),
535 count_bucket(rate * mult, rate * INBOUND_COUNT_SECONDS, rate * INBOUND_COUNT_SECONDS)
541 return !backlog.empty() && size_bucket.
value() > 0 && count_bucket.
value() > 0;
552 bool decrement(
double size)
554 bool size_ok = size_bucket.
decrement(size, -50e3);
555 bool count_ok = count_bucket.
decrement(1, count_floor);
556 return size_ok && count_ok;
563 .count_bucket = count_bucket.
value(),
564 .size_bucket = size_bucket.
value(),
595 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);
597 EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
612 void SetBestBlock(
int height,
std::chrono::seconds time)
override
614 m_best_height = height;
615 m_best_block_time = time;
623 const std::atomic<bool>& interruptMsgProc)
624 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);
636 void ReattemptPrivateBroadcast(
CScheduler& scheduler);
651 void Misbehaving(Peer& peer, const
std::
string& message);
662 bool via_compact_block, const
std::
string& message = "")
671 bool MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer);
684 bool MaybeDisconnectForTxRelayCapacity(
CNode&
node, const
std::
string& msg_type,
699 bool first_time_failure)
724 bool ProcessOrphanTx(Peer& peer)
734 void ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
736 bool via_compact_block)
767 bool IsContinuationOfLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
781 bool TryLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
796 void HeadersDirectFetchBlocks(
CNode& pfrom, const Peer& peer, const
CBlockIndex& last_header);
798 void UpdatePeerStateForReceivedHeaders(
CNode& pfrom, const
CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
805 template <
typename... Args>
806 void MakeAndPushMessage(
CNode&
node, std::string msg_type, Args&&...
args)
const
810 template <
typename... Args>
811 [[maybe_unused]]
void MakeAndPushFeature(
CNode&
node, std::string_view feature_id, Args&&...
args)
const
814 std::vector<unsigned char> feature_data;
821 void PushNodeVersion(
CNode& pnode,
const Peer& peer);
872 std::unique_ptr<TxReconciliationTracker> m_txreconciliation;
875 std::atomic<int> m_best_height{-1};
877 std::atomic<std::chrono::seconds> m_best_block_time{0
s};
885 const Options m_opts;
887 bool RejectIncomingTxs(
const CNode& peer)
const;
895 mutable Mutex m_peer_mutex;
902 std::map<NodeId, PeerRef> m_peer_map
GUARDED_BY(m_peer_mutex);
912 uint32_t GetFetchFlags(
const Peer& peer)
const;
914 std::map<uint64_t, std::chrono::microseconds> m_next_inv_to_inbounds_per_network_key
GUARDED_BY(g_msgproc_mutex);
931 std::atomic<int> m_wtxid_relay_peers{0};
949 std::chrono::microseconds NextInvToInbounds(std::chrono::microseconds now,
950 std::chrono::seconds average_interval,
955 Mutex m_most_recent_block_mutex;
956 std::shared_ptr<const CBlock> m_most_recent_block
GUARDED_BY(m_most_recent_block_mutex);
957 std::shared_ptr<const CBlockHeaderAndShortTxIDs> m_most_recent_compact_block
GUARDED_BY(m_most_recent_block_mutex);
959 std::unique_ptr<const std::map<GenTxid, CTransactionRef>> m_most_recent_block_txs
GUARDED_BY(m_most_recent_block_mutex);
963 Mutex m_headers_presync_mutex;
971 using HeadersPresyncStats = std::pair<arith_uint256, std::optional<std::pair<int64_t, uint32_t>>>;
973 std::map<NodeId, HeadersPresyncStats> m_headers_presync_stats
GUARDED_BY(m_headers_presync_mutex) {};
977 std::atomic_bool m_headers_presync_should_signal{
false};
1047 std::atomic<
std::chrono::seconds> m_last_tip_update{0
s};
1053 void ProcessGetData(
CNode& pfrom, Peer& peer,
const std::atomic<bool>& interruptMsgProc)
1058 void ProcessBlock(
CNode&
node,
const std::shared_ptr<const CBlock>& block,
bool force_processing,
bool min_pow_checked);
1091 std::vector<std::pair<Wtxid, CTransactionRef>> vExtraTxnForCompact
GUARDED_BY(g_msgproc_mutex);
1093 size_t vExtraTxnForCompactIt
GUARDED_BY(g_msgproc_mutex) = 0;
1105 int64_t ApproximateBestBlockDepth() const;
1115 void ProcessGetBlockData(
CNode& pfrom, Peer& peer, const
CInv& inv)
1133 bool PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
1135 const
uint256& stop_hash, uint32_t max_height_diff,
1182 void ProcessAddrs(
std::string_view msg_type,
CNode& pfrom, Peer& peer,
std::vector<
CAddress>&& vAddr, const
std::atomic<
bool>& interruptMsgProc)
1188 void LogBlockHeader(const
CBlockIndex& index, const
CNode& peer,
bool via_compact_block);
1194 InvToSendBucket m_inbound_inv_bucket
GUARDED_BY(m_inv_to_send_mutex);
1195 InvToSendBucket m_outbound_inv_bucket
GUARDED_BY(m_inv_to_send_mutex);
1196 std::atomic<
NodeClock::time_point> m_next_inv_bucket_check{NodeClock::time_point::min()};
1197 std::optional<NodeClock::time_point> m_next_inv_bucket_heartbeat
GUARDED_BY(m_inv_to_send_mutex);
1202const CNodeState* PeerManagerImpl::
State(
NodeId pnode)
const
1204 std::map<NodeId, CNodeState>::const_iterator it = m_node_states.find(pnode);
1205 if (it == m_node_states.end())
1212 return const_cast<CNodeState*
>(std::as_const(*this).State(pnode));
1220static bool IsAddrCompatible(
const Peer& peer,
const CAddress& addr)
1225void PeerManagerImpl::AddAddressKnown(Peer& peer,
const CAddress& addr)
1227 assert(peer.m_addr_known);
1228 peer.m_addr_known->insert(addr.
GetKey());
1231void PeerManagerImpl::PushAddress(Peer& peer,
const CAddress& addr)
1236 assert(peer.m_addr_known);
1237 if (addr.
IsValid() && !peer.m_addr_known->contains(addr.
GetKey()) && IsAddrCompatible(peer, addr)) {
1239 peer.m_addrs_to_send[m_rng.randrange(peer.m_addrs_to_send.size())] = addr;
1241 peer.m_addrs_to_send.push_back(addr);
1246static void AddKnownTx(Peer& peer,
const uint256& hash)
1248 auto tx_relay = peer.GetTxRelay();
1249 if (!tx_relay)
return;
1251 LOCK(tx_relay->m_tx_inventory_mutex);
1252 tx_relay->m_tx_inventory_known_filter.insert(hash);
1256static bool CanServeBlocks(
const Peer& peer)
1263static bool IsLimitedPeer(
const Peer& peer)
1270static bool CanServeWitnesses(
const Peer& peer)
1275std::chrono::microseconds PeerManagerImpl::NextInvToInbounds(std::chrono::microseconds now,
1276 std::chrono::seconds average_interval,
1277 uint64_t network_key)
1279 auto [it, inserted] = m_next_inv_to_inbounds_per_network_key.try_emplace(network_key, 0us);
1280 auto& timer{it->second};
1282 timer = now + m_rng.rand_exp_duration(average_interval);
1287bool PeerManagerImpl::IsBlockRequested(
const uint256& hash)
1289 return mapBlocksInFlight.contains(hash);
1292bool PeerManagerImpl::IsBlockRequestedFromOutbound(
const uint256& hash)
1294 for (
auto range = mapBlocksInFlight.equal_range(hash); range.first != range.second; range.first++) {
1295 auto [nodeid, block_it] = range.first->second;
1296 PeerRef peer{GetPeerRef(nodeid)};
1297 if (peer && !peer->m_is_inbound)
return true;
1303void PeerManagerImpl::RemoveBlockRequest(
const uint256& hash, std::optional<NodeId> from_peer)
1305 auto range = mapBlocksInFlight.equal_range(hash);
1306 if (range.first == range.second) {
1314 while (range.first != range.second) {
1315 const auto& [node_id, list_it]{range.first->second};
1317 if (from_peer && *from_peer != node_id) {
1324 if (state.vBlocksInFlight.begin() == list_it) {
1326 state.m_downloading_since = std::max(state.m_downloading_since, GetTime<std::chrono::microseconds>());
1328 state.vBlocksInFlight.erase(list_it);
1330 if (state.vBlocksInFlight.empty()) {
1332 m_peers_downloading_from--;
1334 state.m_stalling_since = 0us;
1336 range.first = mapBlocksInFlight.erase(range.first);
1340bool PeerManagerImpl::BlockRequested(
NodeId nodeid,
const CBlockIndex& block, std::list<QueuedBlock>::iterator** pit)
1344 CNodeState *state =
State(nodeid);
1345 assert(state !=
nullptr);
1350 for (
auto range = mapBlocksInFlight.equal_range(hash); range.first != range.second; range.first++) {
1351 if (range.first->second.first == nodeid) {
1353 *pit = &range.first->second.second;
1360 RemoveBlockRequest(hash, nodeid);
1362 std::list<QueuedBlock>::iterator it = state->vBlocksInFlight.insert(state->vBlocksInFlight.end(),
1363 {&block, std::unique_ptr<PartiallyDownloadedBlock>(pit ? new PartiallyDownloadedBlock(&m_mempool) : nullptr)});
1364 if (state->vBlocksInFlight.size() == 1) {
1366 state->m_downloading_since = GetTime<std::chrono::microseconds>();
1367 m_peers_downloading_from++;
1369 auto itInFlight = mapBlocksInFlight.insert(std::make_pair(hash, std::make_pair(nodeid, it)));
1371 *pit = &itInFlight->second.second;
1376void PeerManagerImpl::MaybeSetPeerAsAnnouncingHeaderAndIDs(
NodeId nodeid)
1383 if (m_opts.ignore_incoming_txs)
return;
1385 CNodeState* nodestate =
State(nodeid);
1386 PeerRef peer{GetPeerRef(nodeid)};
1387 if (!nodestate || !nodestate->m_provides_cmpctblocks) {
1392 int num_outbound_hb_peers = 0;
1393 for (std::list<NodeId>::iterator it = lNodesAnnouncingHeaderAndIDs.begin(); it != lNodesAnnouncingHeaderAndIDs.end(); it++) {
1394 if (*it == nodeid) {
1395 lNodesAnnouncingHeaderAndIDs.erase(it);
1396 lNodesAnnouncingHeaderAndIDs.push_back(nodeid);
1399 PeerRef peer_ref{GetPeerRef(*it)};
1400 if (peer_ref && !peer_ref->m_is_inbound) ++num_outbound_hb_peers;
1402 if (peer && peer->m_is_inbound) {
1405 if (lNodesAnnouncingHeaderAndIDs.size() >= 3 && num_outbound_hb_peers == 1) {
1406 PeerRef remove_peer{GetPeerRef(lNodesAnnouncingHeaderAndIDs.front())};
1407 if (remove_peer && !remove_peer->m_is_inbound) {
1410 std::swap(lNodesAnnouncingHeaderAndIDs.front(), *std::next(lNodesAnnouncingHeaderAndIDs.begin()));
1419 lNodesAnnouncingHeaderAndIDs.push_back(pfrom->
GetId());
1422 if (nodeid_was_appended && lNodesAnnouncingHeaderAndIDs.size() > 3) {
1425 m_connman.
ForNode(lNodesAnnouncingHeaderAndIDs.front(), [
this](
CNode* pnodeStop) {
1428 pnodeStop->m_bip152_highbandwidth_to =
false;
1431 lNodesAnnouncingHeaderAndIDs.pop_front();
1435bool PeerManagerImpl::TipMayBeStale()
1439 if (m_last_tip_update.load() == 0
s) {
1440 m_last_tip_update = GetTime<std::chrono::seconds>();
1442 return m_last_tip_update.load() < GetTime<std::chrono::seconds>() - std::chrono::seconds{consensusParams.
nPowTargetSpacing * 3} && mapBlocksInFlight.empty();
1445int64_t PeerManagerImpl::ApproximateBestBlockDepth()
const
1450bool PeerManagerImpl::CanDirectFetch()
1457 if (state->pindexBestKnownBlock && pindex == state->pindexBestKnownBlock->GetAncestor(pindex->nHeight))
1459 if (state->pindexBestHeaderSent && pindex == state->pindexBestHeaderSent->GetAncestor(pindex->nHeight))
1464void PeerManagerImpl::ProcessBlockAvailability(
NodeId nodeid) {
1465 CNodeState *state =
State(nodeid);
1466 assert(state !=
nullptr);
1468 if (!state->hashLastUnknownBlock.IsNull()) {
1471 if (state->pindexBestKnownBlock ==
nullptr || pindex->
nChainWork >= state->pindexBestKnownBlock->nChainWork) {
1472 state->pindexBestKnownBlock = pindex;
1474 state->hashLastUnknownBlock.SetNull();
1479void PeerManagerImpl::UpdateBlockAvailability(
NodeId nodeid,
const uint256 &hash) {
1480 CNodeState *state =
State(nodeid);
1481 assert(state !=
nullptr);
1483 ProcessBlockAvailability(nodeid);
1488 if (state->pindexBestKnownBlock ==
nullptr || pindex->
nChainWork >= state->pindexBestKnownBlock->nChainWork) {
1489 state->pindexBestKnownBlock = pindex;
1493 state->hashLastUnknownBlock = hash;
1498void PeerManagerImpl::FindNextBlocksToDownload(
const Peer& peer,
unsigned int count, std::vector<const CBlockIndex*>& vBlocks,
NodeId& nodeStaller)
1503 vBlocks.reserve(vBlocks.size() +
count);
1504 CNodeState *state =
State(peer.m_id);
1505 assert(state !=
nullptr);
1508 ProcessBlockAvailability(peer.m_id);
1510 if (state->pindexBestKnownBlock ==
nullptr || state->pindexBestKnownBlock->nChainWork < m_chainman.
ActiveChain().
Tip()->
nChainWork || state->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
1521 state->pindexBestKnownBlock->GetAncestor(snap_base->nHeight) != snap_base) {
1522 LogDebug(
BCLog::NET,
"Not downloading blocks from peer=%d, which doesn't have the snapshot block in its best chain.\n", peer.m_id);
1531 if (state->pindexLastCommonBlock ==
nullptr ||
1532 fork_point->nChainWork > state->pindexLastCommonBlock->nChainWork ||
1533 state->pindexBestKnownBlock->GetAncestor(state->pindexLastCommonBlock->nHeight) != state->pindexLastCommonBlock) {
1534 state->pindexLastCommonBlock = fork_point;
1536 if (state->pindexLastCommonBlock == state->pindexBestKnownBlock)
1539 const CBlockIndex *pindexWalk = state->pindexLastCommonBlock;
1545 FindNextBlocks(vBlocks, peer, state, pindexWalk,
count, nWindowEnd, &m_chainman.
ActiveChain(), &nodeStaller);
1548void PeerManagerImpl::TryDownloadingHistoricalBlocks(
const Peer& peer,
unsigned int count, std::vector<const CBlockIndex*>& vBlocks,
const CBlockIndex *from_tip,
const CBlockIndex* target_block)
1553 if (vBlocks.size() >=
count) {
1557 vBlocks.reserve(
count);
1560 if (state->pindexBestKnownBlock ==
nullptr || state->pindexBestKnownBlock->GetAncestor(target_block->
nHeight) != target_block) {
1577void 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)
1579 std::vector<const CBlockIndex*> vToFetch;
1580 int nMaxHeight = std::min<int>(state->pindexBestKnownBlock->nHeight, nWindowEnd + 1);
1581 bool is_limited_peer = IsLimitedPeer(peer);
1583 while (pindexWalk->
nHeight < nMaxHeight) {
1587 int nToFetch = std::min(nMaxHeight - pindexWalk->
nHeight, std::max<int>(
count - vBlocks.size(), 128));
1588 vToFetch.resize(nToFetch);
1589 pindexWalk = state->pindexBestKnownBlock->
GetAncestor(pindexWalk->
nHeight + nToFetch);
1590 vToFetch[nToFetch - 1] = pindexWalk;
1591 for (
unsigned int i = nToFetch - 1; i > 0; i--) {
1592 vToFetch[i - 1] = vToFetch[i]->
pprev;
1612 state->pindexLastCommonBlock = pindex;
1619 if (waitingfor == -1) {
1621 waitingfor = mapBlocksInFlight.lower_bound(pindex->
GetBlockHash())->second.first;
1627 if (pindex->
nHeight > nWindowEnd) {
1629 if (vBlocks.size() == 0 && waitingfor != peer.m_id) {
1631 if (nodeStaller) *nodeStaller = waitingfor;
1641 vBlocks.push_back(pindex);
1642 if (vBlocks.size() ==
count) {
1651void PeerManagerImpl::PushNodeVersion(
CNode& pnode,
const Peer& peer)
1653 uint64_t my_services;
1655 uint64_t your_services;
1657 std::string my_user_agent;
1665 my_user_agent =
"/pynode:0.0.1/";
1667 my_tx_relay =
false;
1670 my_services = peer.m_our_services;
1671 my_time = TicksSinceEpoch<std::chrono::seconds>(
NodeClock::now());
1675 my_height = m_best_height;
1676 my_tx_relay = !RejectIncomingTxs(pnode);
1695 BCLog::NET,
"send version message: version=%d, blocks=%d%s, txrelay=%d, peer=%d\n",
1698 my_tx_relay, pnode.
GetId());
1705 if (state) state->m_last_block_announcement = time;
1713 m_node_states.try_emplace(m_node_states.end(), nodeid);
1715 WITH_LOCK(m_tx_download_mutex, m_txdownloadman.CheckIsEmpty(nodeid));
1721 PeerRef peer = std::make_shared<Peer>(nodeid, our_services,
node.IsInboundConn());
1724 m_peer_map.emplace_hint(m_peer_map.end(), nodeid, peer);
1728void PeerManagerImpl::ReattemptInitialBroadcast(
CScheduler& scheduler)
1732 for (
const auto& txid : unbroadcast_txids) {
1735 if (tx !=
nullptr) {
1736 InitiateTxBroadcastToAll(tx->GetWitnessHash());
1745 scheduler.
scheduleFromNow([&] { ReattemptInitialBroadcast(scheduler); }, delta);
1748void PeerManagerImpl::ReattemptPrivateBroadcast(
CScheduler& scheduler)
1752 size_t num_for_rebroadcast{0};
1753 const auto stale_txs = m_tx_for_private_broadcast.GetStale();
1754 if (!stale_txs.empty()) {
1755 for (
const auto& stale_tx : stale_txs) {
1761 "Reattempting broadcast of stale txid=%s wtxid=%s",
1762 stale_tx->GetHash().ToString(), stale_tx->GetWitnessHash().ToString());
1763 ++num_for_rebroadcast;
1766 stale_tx->GetHash().ToString(), stale_tx->GetWitnessHash().ToString(),
1767 mempool_acceptable.m_state.ToString());
1768 m_tx_for_private_broadcast.Remove(stale_tx);
1777 scheduler.
scheduleFromNow([&] { ReattemptPrivateBroadcast(scheduler); }, delta);
1780void PeerManagerImpl::FinalizeNode(
const CNode&
node)
1791 PeerRef peer = RemovePeer(nodeid);
1793 m_wtxid_relay_peers -= peer->m_wtxid_relay;
1794 assert(m_wtxid_relay_peers >= 0);
1796 CNodeState *state =
State(nodeid);
1797 assert(state !=
nullptr);
1799 if (state->fSyncStarted)
1802 for (
const QueuedBlock& entry : state->vBlocksInFlight) {
1803 auto range = mapBlocksInFlight.equal_range(entry.pindex->GetBlockHash());
1804 while (range.first != range.second) {
1805 auto [node_id, list_it] = range.first->second;
1806 if (node_id != nodeid) {
1809 range.first = mapBlocksInFlight.erase(range.first);
1814 LOCK(m_tx_download_mutex);
1815 m_txdownloadman.DisconnectedPeer(nodeid);
1817 if (m_txreconciliation) m_txreconciliation->ForgetPeer(nodeid);
1818 m_num_preferred_download_peers -= state->fPreferredDownload;
1819 m_peers_downloading_from -= (!state->vBlocksInFlight.empty());
1820 assert(m_peers_downloading_from >= 0);
1821 m_outbound_peers_with_protect_from_disconnect -= state->m_chain_sync.m_protect;
1822 assert(m_outbound_peers_with_protect_from_disconnect >= 0);
1824 m_node_states.erase(nodeid);
1826 if (m_node_states.empty()) {
1828 assert(mapBlocksInFlight.empty());
1829 assert(m_num_preferred_download_peers == 0);
1830 assert(m_peers_downloading_from == 0);
1831 assert(m_outbound_peers_with_protect_from_disconnect == 0);
1832 assert(m_wtxid_relay_peers == 0);
1833 WITH_LOCK(m_tx_download_mutex, m_txdownloadman.CheckIsEmpty());
1836 if (
node.fSuccessfullyConnected &&
1837 !
node.IsBlockOnlyConn() && !
node.IsPrivateBroadcastConn() && !
node.IsInboundConn()) {
1845 LOCK(m_headers_presync_mutex);
1846 m_headers_presync_stats.erase(nodeid);
1848 if (
node.IsPrivateBroadcastConn() &&
1849 !m_tx_for_private_broadcast.DidNodeConfirmReception(nodeid) &&
1850 m_tx_for_private_broadcast.HavePendingTransactions()) {
1857bool PeerManagerImpl::HasAllDesirableServiceFlags(
ServiceFlags services)
const
1860 return !(GetDesirableServiceFlags(services) & (~services));
1874PeerRef PeerManagerImpl::GetPeerRef(
NodeId id)
const
1877 auto it = m_peer_map.find(
id);
1878 return it != m_peer_map.end() ? it->second :
nullptr;
1881PeerRef PeerManagerImpl::RemovePeer(
NodeId id)
1885 auto it = m_peer_map.find(
id);
1886 if (it != m_peer_map.end()) {
1887 ret = std::move(it->second);
1888 m_peer_map.erase(it);
1893std::vector<PeerRef> PeerManagerImpl::GetAllPeers()
const
1895 std::vector<PeerRef> peers;
1897 peers.reserve(m_peer_map.size());
1898 for (
const auto& [
_, peer] : m_peer_map) {
1899 peers.push_back(peer);
1908 const CNodeState* state =
State(nodeid);
1909 if (state ==
nullptr)
1911 stats.
nSyncHeight = state->pindexBestKnownBlock ? state->pindexBestKnownBlock->nHeight : -1;
1912 stats.
nCommonHeight = state->pindexLastCommonBlock ? state->pindexLastCommonBlock->nHeight : -1;
1913 for (
const QueuedBlock& queue : state->vBlocksInFlight) {
1920 PeerRef peer = GetPeerRef(nodeid);
1921 if (peer ==
nullptr)
return false;
1929 NodeClock::duration ping_wait{0us};
1930 if ((0 != peer->m_ping_nonce_sent) && (peer->m_ping_start.load() >
NodeClock::epoch)) {
1934 if (
auto tx_relay = peer->GetTxRelay(); tx_relay !=
nullptr) {
1937 LOCK(tx_relay->m_tx_inventory_mutex);
1939 stats.
m_inv_to_send = tx_relay->m_tx_inventory_to_send.size();
1951 LOCK(peer->m_headers_sync_mutex);
1952 if (peer->m_headers_sync) {
1961std::vector<node::TxOrphanage::OrphanInfo> PeerManagerImpl::GetOrphanTransactions()
1963 LOCK(m_tx_download_mutex);
1964 return m_txdownloadman.GetOrphanTransactions();
1969 LOCK(m_inv_to_send_mutex);
1972 .ignores_incoming_txs = m_opts.ignore_incoming_txs,
1973 .private_broadcast = m_opts.private_broadcast,
1974 .tx_send_rate = m_opts.tx_send_rate,
1975 .inbound_bucket = m_inbound_inv_bucket.info(),
1976 .outbound_bucket = m_outbound_inv_bucket.info(),
1980std::vector<PrivateBroadcast::TxBroadcastInfo> PeerManagerImpl::GetPrivateBroadcastInfo()
const
1982 return m_tx_for_private_broadcast.GetBroadcastInfo();
1985std::vector<CTransactionRef> PeerManagerImpl::AbortPrivateBroadcast(
const uint256&
id)
1987 const auto snapshot{m_tx_for_private_broadcast.GetBroadcastInfo()};
1988 std::vector<CTransactionRef> removed_txs;
1990 size_t connections_cancelled{0};
1991 for (
const auto& tx_info : snapshot) {
1993 if (tx->GetHash().ToUint256() !=
id && tx->GetWitnessHash().ToUint256() !=
id)
continue;
1994 if (
const auto peer_acks{m_tx_for_private_broadcast.Remove(tx)}) {
1995 removed_txs.push_back(tx);
2006void PeerManagerImpl::AddToCompactExtraTransactions(
const CTransactionRef& tx)
2008 if (m_opts.max_extra_txs == 0)
return;
2009 if (vExtraTxnForCompact.size() < m_opts.max_extra_txs) {
2010 if (vExtraTxnForCompact.empty()) vExtraTxnForCompact.reserve(m_opts.max_extra_txs);
2011 vExtraTxnForCompact.emplace_back(tx->GetWitnessHash(), tx);
2013 vExtraTxnForCompact[vExtraTxnForCompactIt] = std::make_pair(tx->GetWitnessHash(), tx);
2015 vExtraTxnForCompactIt = (vExtraTxnForCompactIt + 1) % m_opts.max_extra_txs;
2018void PeerManagerImpl::Misbehaving(Peer& peer,
const std::string& message)
2020 LOCK(peer.m_misbehavior_mutex);
2022 const std::string message_prefixed = message.empty() ?
"" : (
": " + message);
2023 peer.m_should_discourage =
true;
2032 bool via_compact_block,
const std::string& message)
2034 PeerRef peer{GetPeerRef(nodeid)};
2045 if (!via_compact_block) {
2046 if (peer) Misbehaving(*peer, message);
2054 if (peer && !via_compact_block && !peer->m_is_inbound) {
2055 if (peer) Misbehaving(*peer, message);
2062 if (peer) Misbehaving(*peer, message);
2066 if (peer) Misbehaving(*peer, message);
2071 if (message !=
"") {
2076bool PeerManagerImpl::BlockRequestAllowed(
const CBlockIndex& block_index)
2096 PeerRef peer = GetPeerRef(peer_id);
2103 RemoveBlockRequest(block_index.
GetBlockHash(), std::nullopt);
2106 if (!BlockRequested(peer_id, block_index))
return util::Unexpected{
"Already requested from this peer"};
2129 return std::make_unique<PeerManagerImpl>(connman, addrman, banman, chainman, pool, warnings, opts);
2135 : m_rng{opts.deterministic_rng},
2137 m_chainparams(chainman.GetParams()),
2141 m_chainman(chainman),
2143 m_txdownloadman{
node::TxDownloadOptions{pool, opts.deterministic_rng}},
2144 m_warnings{warnings},
2146 m_inbound_inv_bucket(m_opts.tx_send_rate, 1.0),
2151 if (opts.reconcile_txs) {
2156void PeerManagerImpl::StartScheduledTasks(
CScheduler& scheduler)
2167 scheduler.
scheduleFromNow([&] { ReattemptInitialBroadcast(scheduler); }, delta);
2169 if (m_opts.private_broadcast) {
2170 scheduler.
scheduleFromNow([&] { ReattemptPrivateBroadcast(scheduler); }, 0min);
2174void PeerManagerImpl::ActiveTipChange(
const CBlockIndex& new_tip,
bool is_ibd)
2182 LOCK(m_tx_download_mutex);
2186 m_txdownloadman.ActiveTipChange();
2196void PeerManagerImpl::BlockConnected(
2198 const std::shared_ptr<const CBlock>& pblock,
2203 m_last_tip_update = GetTime<std::chrono::seconds>();
2206 auto stalling_timeout = m_block_stalling_timeout.load();
2210 if (m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
2219 LOCK(m_tx_download_mutex);
2220 m_txdownloadman.BlockConnected(pblock);
2224void PeerManagerImpl::BlockDisconnected(
const std::shared_ptr<const CBlock> &block,
const CBlockIndex* pindex)
2226 LOCK(m_tx_download_mutex);
2227 m_txdownloadman.BlockDisconnected();
2234void PeerManagerImpl::NewPoWValidBlock(
const CBlockIndex *pindex,
const std::shared_ptr<const CBlock>& pblock)
2236 auto pcmpctblock = std::make_shared<const CBlockHeaderAndShortTxIDs>(*pblock,
FastRandomContext().rand64());
2240 if (pindex->
nHeight <= m_highest_fast_announce)
2242 m_highest_fast_announce = pindex->
nHeight;
2246 uint256 hashBlock(pblock->GetHash());
2247 const std::shared_future<CSerializedNetMsg> lazy_ser{
2251 auto most_recent_block_txs = std::make_unique<std::map<GenTxid, CTransactionRef>>();
2252 for (
const auto& tx : pblock->vtx) {
2253 most_recent_block_txs->emplace(tx->GetHash(), tx);
2254 most_recent_block_txs->emplace(tx->GetWitnessHash(), tx);
2257 LOCK(m_most_recent_block_mutex);
2258 m_most_recent_block_hash = hashBlock;
2259 m_most_recent_block = pblock;
2260 m_most_recent_compact_block = pcmpctblock;
2261 m_most_recent_block_txs = std::move(most_recent_block_txs);
2269 ProcessBlockAvailability(pnode->
GetId());
2273 if (state.m_requested_hb_cmpctblocks && !PeerHasHeader(&state, pindex) && PeerHasHeader(&state, pindex->
pprev)) {
2275 LogDebug(
BCLog::NET,
"%s sending header-and-ids %s to peer=%d\n",
"PeerManager::NewPoWValidBlock",
2276 hashBlock.ToString(), pnode->
GetId());
2279 PushMessage(*pnode, ser_cmpctblock.Copy());
2280 state.pindexBestHeaderSent = pindex;
2289void PeerManagerImpl::UpdatedBlockTip(
const CBlockIndex *pindexNew,
const CBlockIndex *pindexFork,
bool fInitialDownload)
2291 SetBestBlock(pindexNew->
nHeight, std::chrono::seconds{pindexNew->GetBlockTime()});
2294 if (fInitialDownload)
return;
2297 std::vector<uint256> vHashes;
2299 while (pindexToAnnounce != pindexFork) {
2301 pindexToAnnounce = pindexToAnnounce->
pprev;
2311 for (
auto& it : m_peer_map) {
2312 Peer& peer = *it.second;
2313 LOCK(peer.m_block_inv_mutex);
2314 for (
const uint256& hash : vHashes | std::views::reverse) {
2315 peer.m_blocks_for_headers_relay.push_back(hash);
2327void PeerManagerImpl::BlockChecked(
const std::shared_ptr<const CBlock>& block,
const BlockValidationState& state)
2331 const uint256 hash(block->GetHash());
2332 std::map<uint256, std::pair<NodeId, bool>>::iterator it = mapBlockSource.find(hash);
2337 it != mapBlockSource.end() &&
2338 State(it->second.first)) {
2339 MaybePunishNodeForBlock( it->second.first, state, !it->second.second);
2349 mapBlocksInFlight.count(hash) == mapBlocksInFlight.size()) {
2350 if (it != mapBlockSource.end()) {
2351 MaybeSetPeerAsAnnouncingHeaderAndIDs(it->second.first);
2354 if (it != mapBlockSource.end())
2355 mapBlockSource.erase(it);
2363bool PeerManagerImpl::AlreadyHaveBlock(
const uint256& block_hash)
2368void PeerManagerImpl::SendPings()
2371 for(
auto& it : m_peer_map) it.second->m_ping_queued =
true;
2374std::vector<Wtxid> InvToSendBucket::TakeForProcessing(
CTxMemPool& mempool)
2378 size_t n_to_take =
static_cast<size_t>(std::max<double>(count_bucket.
value() - count_floor, 0));
2380 std::vector<Wtxid> best;
2383 bool tokens_left =
true;
2384 for (
auto txiter : itervec) {
2385 auto& wtxid = txiter->GetTx().GetWitnessHash();
2387 best.push_back(wtxid);
2388 if (!decrement(txiter->GetTx().ComputeTotalSize())) {
2389 tokens_left =
false;
2392 backlog.push_back(wtxid);
2398 std::vector<Wtxid> dummy;
2400 dummy.swap(backlog);
2410 if (!backlog_bumped && now <= m_next_inv_bucket_check.load())
return;
2413 LOCK(m_inv_to_send_mutex);
2414 m_inbound_inv_bucket.increment(now);
2415 m_outbound_inv_bucket.increment(now);
2418 if (!m_next_inv_bucket_heartbeat.has_value()) {
2420 m_next_inv_bucket_heartbeat = now;
2423 if (m_next_inv_bucket_heartbeat.has_value() && now >= *m_next_inv_bucket_heartbeat) {
2424 LogDebug(
BCLog::NET,
"Transaction rate-limiting backlog inbound=%d itok=%.1f isz=%.1f outbound=%d otok=%.1f osz=%.1f",
2425 m_inbound_inv_bucket.backlog.size(),
2426 m_inbound_inv_bucket.count_bucket.value(),
2427 m_inbound_inv_bucket.size_bucket.value(),
2428 m_outbound_inv_bucket.backlog.size(),
2429 m_outbound_inv_bucket.count_bucket.value(),
2430 m_outbound_inv_bucket.size_bucket.value());
2431 if (m_inbound_inv_bucket.backlog.empty() && m_outbound_inv_bucket.backlog.empty()) {
2432 m_next_inv_bucket_heartbeat = std::nullopt;
2439 bool in_avail = m_inbound_inv_bucket.avail();
2440 bool out_avail = m_outbound_inv_bucket.avail();
2441 if (!in_avail && !out_avail)
return;
2443 std::vector<Wtxid> for_inbound;
2444 std::vector<Wtxid> for_outbound;
2448 if (in_avail) for_inbound = m_inbound_inv_bucket.TakeForProcessing(m_mempool);
2449 if (out_avail) for_outbound = m_outbound_inv_bucket.TakeForProcessing(m_mempool);
2452 if (!for_inbound.empty() || !for_outbound.empty()) {
2453 bool any_inbound_connected =
false;
2454 bool any_outbound_connected =
false;
2455 for (
const PeerRef& peer_ref : GetAllPeers()) {
2456 if (!peer_ref)
continue;
2457 Peer& peer{*peer_ref};
2458 auto tx_relay = peer.GetTxRelay();
2459 if (!tx_relay)
continue;
2461 LOCK(tx_relay->m_tx_inventory_mutex);
2467 if (tx_relay->m_next_inv_send_time == 0
s)
continue;
2468 if (peer.m_is_inbound) {
2469 any_inbound_connected =
true;
2471 any_outbound_connected =
true;
2473 for (
auto& i : (peer.m_is_inbound ? for_inbound : for_outbound)) {
2474 tx_relay->m_tx_inventory_to_send.push_back(i);
2481 if (!any_inbound_connected) m_inbound_inv_bucket.backlog.clear();
2482 if (!any_outbound_connected) m_outbound_inv_bucket.backlog.clear();
2486void PeerManagerImpl::InitiateTxBroadcastToAll(
const Wtxid& wtxid)
2489 LOCK(m_inv_to_send_mutex);
2490 m_inbound_inv_bucket.backlog.push_back(wtxid);
2491 m_outbound_inv_bucket.backlog.push_back(wtxid);
2498 const auto txstr{
strprintf(
"txid=%s, wtxid=%s", tx->GetHash().ToString(), tx->GetWitnessHash().ToString())};
2499 switch (m_tx_for_private_broadcast.Add(tx)) {
2514void PeerManagerImpl::RelayAddress(
NodeId originator,
2530 const auto current_time{GetTime<std::chrono::seconds>()};
2538 unsigned int nRelayNodes = (fReachable || (hasher.Finalize() & 1)) ? 2 : 1;
2540 std::array<std::pair<uint64_t, Peer*>, 2> best{{{0,
nullptr}, {0,
nullptr}}};
2541 assert(nRelayNodes <= best.size());
2545 for (
auto& [
id, peer] : m_peer_map) {
2546 if (peer->m_addr_relay_enabled &&
id != originator && IsAddrCompatible(*peer, addr)) {
2548 for (
unsigned int i = 0; i < nRelayNodes; i++) {
2549 if (hashKey > best[i].first) {
2550 std::copy(best.begin() + i, best.begin() + nRelayNodes - 1, best.begin() + i + 1);
2551 best[i] = std::make_pair(hashKey, peer.get());
2558 for (
unsigned int i = 0; i < nRelayNodes && best[i].first != 0; i++) {
2559 PushAddress(*best[i].second, addr);
2563void PeerManagerImpl::ProcessGetBlockData(
CNode& pfrom, Peer& peer,
const CInv& inv)
2573 std::shared_ptr<const CBlock> a_recent_block;
2574 std::shared_ptr<const CBlockHeaderAndShortTxIDs> a_recent_compact_block;
2576 LOCK(m_most_recent_block_mutex);
2577 a_recent_block = m_most_recent_block;
2578 a_recent_compact_block = m_most_recent_compact_block;
2581 bool need_activate_chain =
false;
2593 need_activate_chain =
true;
2597 if (need_activate_chain) {
2599 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
2606 bool can_direct_fetch{
false};
2614 if (!BlockRequestAllowed(*pindex)) {
2615 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block that isn't in the main chain\n", __func__, pfrom.
GetId());
2642 can_direct_fetch = CanDirectFetch();
2646 std::shared_ptr<const CBlock> pblock;
2647 if (a_recent_block && a_recent_block->GetHash() == inv.
hash) {
2648 pblock = a_recent_block;
2666 std::shared_ptr<CBlock> pblockRead = std::make_shared<CBlock>();
2676 pblock = pblockRead;
2684 bool sendMerkleBlock =
false;
2686 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
2687 LOCK(tx_relay->m_bloom_filter_mutex);
2688 if (tx_relay->m_bloom_filter) {
2689 sendMerkleBlock =
true;
2690 merkleBlock =
CMerkleBlock(*pblock, *tx_relay->m_bloom_filter);
2693 if (sendMerkleBlock) {
2701 for (
const auto& [tx_idx,
_] : merkleBlock.
vMatchedTxn)
2712 if (a_recent_compact_block && a_recent_compact_block->header.GetHash() == inv.
hash) {
2725 LOCK(peer.m_block_inv_mutex);
2727 if (inv.
hash == peer.m_continuation_block) {
2731 std::vector<CInv> vInv;
2732 vInv.emplace_back(
MSG_BLOCK, tip->GetBlockHash());
2734 peer.m_continuation_block.SetNull();
2742 auto txinfo{std::visit(
2743 [&](
const auto&
id) {
2744 return m_mempool.
info_for_relay(
id,
WITH_LOCK(tx_relay.m_tx_inventory_mutex,
return tx_relay.m_last_inv_sequence));
2748 return std::move(txinfo.tx);
2753 LOCK(m_most_recent_block_mutex);
2754 if (m_most_recent_block_txs !=
nullptr) {
2755 auto it = m_most_recent_block_txs->find(gtxid);
2756 if (it != m_most_recent_block_txs->end())
return it->second;
2763void PeerManagerImpl::ProcessGetData(
CNode& pfrom, Peer& peer,
const std::atomic<bool>& interruptMsgProc)
2767 auto tx_relay = peer.GetTxRelay();
2769 std::deque<CInv>::iterator it = peer.m_getdata_requests.begin();
2770 std::vector<CInv> vNotFound;
2775 while (it != peer.m_getdata_requests.end() && it->IsGenTxMsg()) {
2776 if (interruptMsgProc)
return;
2781 const CInv &inv = *it++;
2783 if (tx_relay ==
nullptr) {
2789 if (
auto tx{FindTxForGetData(*tx_relay,
ToGenTxid(inv))}) {
2792 MakeAndPushMessage(pfrom,
NetMsgType::TX, maybe_with_witness(*tx));
2795 vNotFound.push_back(inv);
2801 if (it != peer.m_getdata_requests.end() && !pfrom.
fPauseSend) {
2802 const CInv &inv = *it++;
2804 ProcessGetBlockData(pfrom, peer, inv);
2813 peer.m_getdata_requests.erase(peer.m_getdata_requests.begin(), it);
2815 if (!vNotFound.empty()) {
2834uint32_t PeerManagerImpl::GetFetchFlags(
const Peer& peer)
const
2836 uint32_t nFetchFlags = 0;
2837 if (CanServeWitnesses(peer)) {
2846 for (
size_t i = 0; i < req.
indexes.size(); i++) {
2848 Misbehaving(peer,
"getblocktxn with out-of-bounds tx indices");
2855 uint32_t tx_requested_size{0};
2856 for (
const auto& tx : resp.txn) tx_requested_size += tx->ComputeTotalSize();
2862bool PeerManagerImpl::CheckHeadersPoW(
const std::vector<CBlockHeader>&
headers, Peer& peer)
2866 Misbehaving(peer,
"header with invalid proof of work");
2871 if (!CheckHeadersAreContinuous(
headers)) {
2872 Misbehaving(peer,
"non-continuous headers sequence");
2897void PeerManagerImpl::HandleUnconnectingHeaders(
CNode& pfrom, Peer& peer,
2898 const std::vector<CBlockHeader>&
headers)
2902 if (MaybeSendGetHeaders(pfrom,
GetLocator(best_header), peer)) {
2903 LogDebug(
BCLog::NET,
"received header %s: missing prev block %s, sending getheaders (%d) to end (peer=%d)\n",
2905 headers[0].hashPrevBlock.ToString(),
2906 best_header->nHeight,
2916bool PeerManagerImpl::CheckHeadersAreContinuous(
const std::vector<CBlockHeader>&
headers)
const
2920 if (!hashLastBlock.
IsNull() && header.hashPrevBlock != hashLastBlock) {
2923 hashLastBlock = header.GetHash();
2928bool PeerManagerImpl::IsContinuationOfLowWorkHeadersSync(Peer& peer,
CNode& pfrom, std::vector<CBlockHeader>&
headers)
2930 if (peer.m_headers_sync) {
2931 auto result = peer.m_headers_sync->ProcessNextHeaders(
headers,
headers.size() == m_opts.max_headers_result);
2933 if (result.success) peer.m_last_getheaders_timestamp = {};
2934 if (result.request_more) {
2935 auto locator = peer.m_headers_sync->NextHeadersRequestLocator();
2937 Assume(!locator.vHave.empty());
2940 if (!locator.vHave.empty()) {
2943 bool sent_getheaders = MaybeSendGetHeaders(pfrom, locator, peer);
2946 locator.vHave.front().ToString(), pfrom.
GetId());
2951 peer.m_headers_sync.reset(
nullptr);
2956 LOCK(m_headers_presync_mutex);
2957 m_headers_presync_stats.erase(pfrom.
GetId());
2960 HeadersPresyncStats stats;
2961 stats.first = peer.m_headers_sync->GetPresyncWork();
2963 stats.second = {peer.m_headers_sync->GetPresyncHeight(),
2964 peer.m_headers_sync->GetPresyncTime()};
2968 LOCK(m_headers_presync_mutex);
2969 m_headers_presync_stats[pfrom.
GetId()] = stats;
2970 auto best_it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
2971 bool best_updated =
false;
2972 if (best_it == m_headers_presync_stats.end()) {
2976 const HeadersPresyncStats* stat_best{
nullptr};
2977 for (
const auto& [peer, stat] : m_headers_presync_stats) {
2978 if (!stat_best || stat > *stat_best) {
2983 m_headers_presync_bestpeer = peer_best;
2984 best_updated = (peer_best == pfrom.
GetId());
2985 }
else if (best_it->first == pfrom.
GetId() || stats > best_it->second) {
2987 m_headers_presync_bestpeer = pfrom.
GetId();
2988 best_updated =
true;
2990 if (best_updated && stats.second.has_value()) {
2992 m_headers_presync_should_signal =
true;
2996 if (result.success) {
2999 headers.swap(result.pow_validated_headers);
3002 return result.success;
3010bool PeerManagerImpl::TryLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
const CBlockIndex& chain_start_header, std::vector<CBlockHeader>&
headers)
3017 arith_uint256 minimum_chain_work = GetAntiDoSWorkThreshold();
3021 if (total_work < minimum_chain_work) {
3025 if (
headers.size() == m_opts.max_headers_result) {
3035 LOCK(peer.m_headers_sync_mutex);
3038 m_chainparams.
HeadersSync(), chain_start_header, minimum_chain_work));
3044 const auto msg{
strprintf(
"Failure when attempting to initiate headers sync: %s", e.what())};
3045 std::cerr <<
msg << std::endl;
3053 (void)IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3067bool PeerManagerImpl::IsAncestorOfBestHeaderOrTip(
const CBlockIndex* header)
3069 if (header ==
nullptr) {
3071 }
else if (m_chainman.m_best_header !=
nullptr && header == m_chainman.m_best_header->GetAncestor(header->
nHeight)) {
3079bool PeerManagerImpl::MaybeSendGetHeaders(
CNode& pfrom,
const CBlockLocator& locator, Peer& peer)
3087 peer.m_last_getheaders_timestamp = current_time;
3098void PeerManagerImpl::HeadersDirectFetchBlocks(
CNode& pfrom,
const Peer& peer,
const CBlockIndex& last_header)
3101 CNodeState *nodestate =
State(pfrom.
GetId());
3104 std::vector<const CBlockIndex*> vToFetch;
3112 vToFetch.push_back(pindexWalk);
3114 pindexWalk = pindexWalk->
pprev;
3126 std::vector<CInv> vGetData;
3128 for (
const CBlockIndex* pindex : vToFetch | std::views::reverse) {
3133 uint32_t nFetchFlags = GetFetchFlags(peer);
3135 BlockRequested(pfrom.
GetId(), *pindex);
3139 if (vGetData.size() > 1) {
3144 if (vGetData.size() > 0) {
3145 if (!m_opts.ignore_incoming_txs &&
3146 nodestate->m_provides_cmpctblocks &&
3147 vGetData.size() == 1 &&
3148 mapBlocksInFlight.size() == 1 &&
3164void PeerManagerImpl::UpdatePeerStateForReceivedHeaders(
CNode& pfrom,
3165 const CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
3168 CNodeState *nodestate =
State(pfrom.
GetId());
3185 if (nodestate->pindexBestKnownBlock && nodestate->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
3207 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) {
3209 nodestate->m_chain_sync.m_protect =
true;
3210 ++m_outbound_peers_with_protect_from_disconnect;
3215void PeerManagerImpl::ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
3216 std::vector<CBlockHeader>&&
headers,
3217 bool via_compact_block)
3219 size_t nCount =
headers.size();
3226 LOCK(peer.m_headers_sync_mutex);
3227 if (peer.m_headers_sync) {
3228 peer.m_headers_sync.reset(
nullptr);
3229 LOCK(m_headers_presync_mutex);
3230 m_headers_presync_stats.erase(pfrom.
GetId());
3234 peer.m_last_getheaders_timestamp = {};
3242 if (!CheckHeadersPoW(
headers, peer)) {
3257 bool already_validated_work =
false;
3260 bool have_headers_sync =
false;
3262 LOCK(peer.m_headers_sync_mutex);
3264 already_validated_work = IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3280 have_headers_sync = !!peer.m_headers_sync;
3285 bool headers_connect_blockindex{chain_start_header !=
nullptr};
3287 if (!headers_connect_blockindex) {
3291 HandleUnconnectingHeaders(pfrom, peer,
headers);
3298 peer.m_last_getheaders_timestamp = {};
3308 already_validated_work = already_validated_work || IsAncestorOfBestHeaderOrTip(last_received_header);
3315 already_validated_work =
true;
3321 if (!already_validated_work && TryLowWorkHeadersSync(peer, pfrom,
3322 *chain_start_header,
headers)) {
3334 bool received_new_header{last_received_header ==
nullptr};
3340 state, &pindexLast)};
3346 "If this happens with all peers, consider database corruption (that -reindex may fix) "
3347 "or a potential consensus incompatibility.",
3350 MaybePunishNodeForBlock(pfrom.
GetId(), state, via_compact_block,
"invalid header received");
3356 if (processed && received_new_header) {
3357 LogBlockHeader(*pindexLast, pfrom,
false);
3361 if (nCount == m_opts.max_headers_result && !have_headers_sync) {
3363 if (MaybeSendGetHeaders(pfrom,
GetLocator(pindexLast), peer)) {
3368 UpdatePeerStateForReceivedHeaders(pfrom, *pindexLast, received_new_header, nCount == m_opts.max_headers_result);
3371 HeadersDirectFetchBlocks(pfrom, peer, *pindexLast);
3377 bool first_time_failure)
3383 PeerRef peer{GetPeerRef(nodeid)};
3386 ptx->GetHash().ToString(),
3387 ptx->GetWitnessHash().ToString(),
3391 const auto& [add_extra_compact_tx, unique_parents, package_to_validate] = m_txdownloadman.MempoolRejectedTx(ptx, state, nodeid, first_time_failure);
3394 AddToCompactExtraTransactions(ptx);
3396 for (
const Txid& parent_txid : unique_parents) {
3397 if (peer) AddKnownTx(*peer, parent_txid.ToUint256());
3400 return package_to_validate;
3403void PeerManagerImpl::ProcessValidTx(
NodeId nodeid,
const CTransactionRef& tx,
const std::list<CTransactionRef>& replaced_transactions)
3409 m_txdownloadman.MempoolAcceptedTx(tx);
3413 tx->GetHash().ToString(),
3414 tx->GetWitnessHash().ToString(),
3417 InitiateTxBroadcastToAll(tx->GetWitnessHash());
3420 AddToCompactExtraTransactions(removedTx);
3430 const auto&
package = package_to_validate.m_txns;
3431 const auto& senders = package_to_validate.
m_senders;
3434 m_txdownloadman.MempoolRejectedPackage(package);
3437 if (!
Assume(package.size() == 2))
return;
3441 auto package_iter = package.rbegin();
3442 auto senders_iter = senders.rbegin();
3443 while (package_iter != package.rend()) {
3444 const auto& tx = *package_iter;
3445 const NodeId nodeid = *senders_iter;
3446 const auto it_result{package_result.
m_tx_results.find(tx->GetWitnessHash())};
3450 const auto& tx_result = it_result->second;
3451 switch (tx_result.m_result_type) {
3454 ProcessValidTx(nodeid, tx, tx_result.m_replaced_transactions);
3464 ProcessInvalidTx(nodeid, tx, tx_result.m_state,
false);
3482bool PeerManagerImpl::ProcessOrphanTx(Peer& peer)
3487 while (
CTransactionRef porphanTx = m_txdownloadman.GetTxToReconsider(peer.m_id)) {
3490 const Txid& orphanHash = porphanTx->GetHash();
3491 const Wtxid& orphan_wtxid = porphanTx->GetWitnessHash();
3508 ProcessInvalidTx(peer.m_id, porphanTx, state,
false);
3517bool PeerManagerImpl::PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
3519 const uint256& stop_hash, uint32_t max_height_diff,
3523 const bool supported_filter_type =
3526 if (!supported_filter_type) {
3528 static_cast<uint8_t
>(filter_type),
node.DisconnectMsg());
3529 node.fDisconnect =
true;
3538 if (!stop_index || !BlockRequestAllowed(*stop_index)) {
3541 node.fDisconnect =
true;
3546 uint32_t stop_height = stop_index->
nHeight;
3547 if (start_height > stop_height) {
3549 "start height %d and stop height %d, %s",
3550 start_height, stop_height,
node.DisconnectMsg());
3551 node.fDisconnect =
true;
3554 if (stop_height - start_height >= max_height_diff) {
3556 stop_height - start_height + 1, max_height_diff,
node.DisconnectMsg());
3557 node.fDisconnect =
true;
3562 if (!filter_index) {
3572 uint8_t filter_type_ser;
3573 uint32_t start_height;
3576 vRecv >> filter_type_ser >> start_height >> stop_hash;
3582 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3587 std::vector<BlockFilter> filters;
3589 LogDebug(
BCLog::NET,
"Failed to find block filter in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3594 for (
const auto& filter : filters) {
3601 uint8_t filter_type_ser;
3602 uint32_t start_height;
3605 vRecv >> filter_type_ser >> start_height >> stop_hash;
3611 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3617 if (start_height > 0) {
3619 stop_index->
GetAncestor(
static_cast<int>(start_height - 1));
3621 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3627 std::vector<uint256> filter_hashes;
3629 LogDebug(
BCLog::NET,
"Failed to find block filter hashes in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3643 uint8_t filter_type_ser;
3646 vRecv >> filter_type_ser >> stop_hash;
3652 if (!PrepareBlockFilterRequest(
node, peer, filter_type, 0, stop_hash,
3653 std::numeric_limits<uint32_t>::max(),
3654 stop_index, filter_index)) {
3662 for (
int i =
int(
headers.size()) - 1; i >= 0; i--) {
3667 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3681 bool new_block{
false};
3682 m_chainman.
ProcessNewBlock(block, force_processing, min_pow_checked, &new_block);
3684 node.m_last_block_time = GetTime<std::chrono::seconds>();
3689 RemoveBlockRequest(block->GetHash(), std::nullopt);
3692 mapBlockSource.erase(block->GetHash());
3696void PeerManagerImpl::ProcessCompactBlockTxns(
CNode& pfrom, Peer& peer,
const BlockTransactions& block_transactions)
3698 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
3699 bool fBlockRead{
false};
3703 auto range_flight = mapBlocksInFlight.equal_range(block_transactions.
blockhash);
3704 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
3705 bool requested_block_from_this_peer{
false};
3708 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
3710 while (range_flight.first != range_flight.second) {
3711 auto [node_id, block_it] = range_flight.first->second;
3712 if (node_id == pfrom.
GetId() && block_it->partialBlock) {
3713 requested_block_from_this_peer =
true;
3716 range_flight.first++;
3719 if (!requested_block_from_this_peer) {
3731 Misbehaving(peer,
"previous compact block reconstruction attempt failed");
3742 Misbehaving(peer,
"invalid compact block/non-matching block transactions");
3745 if (first_in_flight) {
3750 std::vector<CInv> invs;
3755 LogDebug(
BCLog::NET,
"Peer %d sent us a compact block but it failed to reconstruct, waiting on first download to complete\n", pfrom.
GetId());
3768 mapBlockSource.emplace(block_transactions.
blockhash, std::make_pair(pfrom.
GetId(),
false));
3783void PeerManagerImpl::LogBlockHeader(
const CBlockIndex& index,
const CNode& peer,
bool via_compact_block) {
3795 "Saw new %sheader hash=%s height=%d %s",
3796 via_compact_block ?
"cmpctblock " :
"",
3808void PeerManagerImpl::PushPrivateBroadcastTx(
CNode&
node)
3812 const auto opt_tx{m_tx_for_private_broadcast.PickTxForSend(
node.GetId(),
CService{node.addr})};
3815 node.fDisconnect =
true;
3821 tx->GetHash().ToString(), tx->HasWitness() ?
strprintf(
", wtxid=%s", tx->GetWitnessHash().ToString()) :
"",
3827void PeerManagerImpl::ProcessMessage(Peer& peer,
CNode& pfrom,
const std::string& msg_type,
DataStream& vRecv,
3829 const std::atomic<bool>& interruptMsgProc)
3844 uint64_t nNonce = 1;
3847 std::string cleanSubVer;
3848 int starting_height = -1;
3851 vRecv >> nVersion >> Using<CustomUintFormatter<8>>(nServices) >> nTime;
3866 LogDebug(
BCLog::NET,
"peer does not offer the expected services (%08x offered, %08x expected), %s",
3868 GetDesirableServiceFlags(nServices),
3881 if (!vRecv.
empty()) {
3889 if (!vRecv.
empty()) {
3890 std::string strSubVer;
3894 if (!vRecv.
empty()) {
3895 vRecv >> starting_height;
3915 PushNodeVersion(pfrom, peer);
3919 const int greatest_common_version = std::min(nVersion, pfrom.
AdvertisedVersion());
3924 peer.m_their_services = nServices;
3928 pfrom.cleanSubVer = cleanSubVer;
3939 (fRelay || (peer.m_our_services &
NODE_BLOOM))) {
3940 auto*
const tx_relay = peer.SetTxRelay();
3942 LOCK(tx_relay->m_bloom_filter_mutex);
3943 tx_relay->m_relay_txs = fRelay;
3949 LogDebug(
BCLog::NET,
"receive version message: %s: version %d, blocks=%d, us=%s, txrelay=%d, %s%s",
3950 cleanSubVer.empty() ?
"<no user agent>" : cleanSubVer, pfrom.
nVersion,
3952 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
3970 if (greatest_common_version >= 70016) {
3985 const auto* tx_relay = peer.GetTxRelay();
3986 if (tx_relay &&
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs) &&
3988 const uint64_t recon_salt = m_txreconciliation->PreRegisterPeer(pfrom.
GetId());
4001 if (MaybeDisconnectForTxRelayCapacity(pfrom, msg_type, pfrom.
GetId()))
return;
4009 m_num_preferred_download_peers += state->fPreferredDownload;
4015 bool send_getaddr{
false};
4017 send_getaddr = SetupAddressRelay(pfrom, peer);
4050 peer.m_time_offset =
NodeSeconds{std::chrono::seconds{nTime}} - Now<NodeSeconds>();
4054 m_outbound_time_offsets.Add(peer.m_time_offset);
4055 m_outbound_time_offsets.WarnIfOutOfSync();
4059 if (greatest_common_version <= 70012) {
4060 constexpr auto finalAlert{
"60010000000000000000000000ffffff7f00000000ffffff7ffeffff7f01ffffff7f00000000ffffff7f00ffffff7f002f555247454e543a20416c657274206b657920636f6d70726f6d697365642c2075706772616465207265717569726564004630440220653febd6410f470f6bae11cad19c48413becb1ac2c17f908fd0fd53bdc3abd5202206d0e9c96fe88d4a0f01ed9dedae2b6f9e00da94cad0fecaae66ecf689bf71b50"_hex};
4061 MakeAndPushMessage(pfrom,
"alert", finalAlert);
4084 auto new_peer_msg = [&]() {
4086 return strprintf(
"New %s peer connected: transport: %s, version: %d, %s%s",
4090 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
4098 LogInfo(
"%s", new_peer_msg());
4101 if (
auto tx_relay = peer.GetTxRelay()) {
4110 tx_relay->m_tx_inventory_mutex,
4111 return tx_relay->m_tx_inventory_to_send.empty() &&
4112 tx_relay->m_next_inv_send_time == 0
s));
4122 PushPrivateBroadcastTx(pfrom);
4135 if (m_txreconciliation) {
4136 if (!peer.m_wtxid_relay || !m_txreconciliation->IsPeerRegistered(pfrom.
GetId())) {
4140 m_txreconciliation->ForgetPeer(pfrom.
GetId());
4146 const CNodeState* state =
State(pfrom.
GetId());
4148 .m_preferred = state->fPreferredDownload,
4149 .m_relay_permissions = pfrom.HasPermission(NetPermissionFlags::Relay),
4150 .m_wtxid_relay = peer.m_wtxid_relay,
4159 peer.m_prefers_headers =
true;
4164 uint8_t sendcmpct_hb{0};
4165 uint64_t sendcmpct_version{0};
4166 vRecv >> sendcmpct_hb >> sendcmpct_version;
4170 if (sendcmpct_hb > 1) {
4171 Misbehaving(peer,
"invalid sendcmpct announce field");
4179 CNodeState* nodestate =
State(pfrom.
GetId());
4180 nodestate->m_provides_cmpctblocks =
true;
4181 nodestate->m_requested_hb_cmpctblocks = sendcmpct_hb;
4198 if (!peer.m_wtxid_relay) {
4199 peer.m_wtxid_relay =
true;
4200 m_wtxid_relay_peers++;
4219 peer.m_wants_addrv2 =
true;
4236 std::string feature_id;
4240 std::vector<unsigned char> feature_data_vec;
4243 }
catch (
const std::exception&) {
4246 if (feature_id.size() < 4 || !vRecv.
empty()) {
4266 if (!m_txreconciliation) {
4267 LogDebug(
BCLog::NET,
"sendtxrcncl from peer=%d ignored, as our node does not have txreconciliation enabled\n", pfrom.
GetId());
4278 if (RejectIncomingTxs(pfrom)) {
4287 const auto* tx_relay = peer.GetTxRelay();
4288 if (!tx_relay || !
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs)) {
4294 uint32_t peer_txreconcl_version;
4295 uint64_t remote_salt;
4296 vRecv >> peer_txreconcl_version >> remote_salt;
4299 peer_txreconcl_version, remote_salt);
4331 const auto ser_params{
4339 std::vector<CAddress> vAddr;
4340 vRecv >> ser_params(vAddr);
4341 ProcessAddrs(msg_type, pfrom, peer, std::move(vAddr), interruptMsgProc);
4346 std::vector<CInv> vInv;
4350 Misbehaving(peer,
strprintf(
"inv message size = %u", vInv.size()));
4354 const bool reject_tx_invs{RejectIncomingTxs(pfrom)};
4355 std::unordered_set<uint256, SaltedUint256Hasher> seen_txids{0, m_txhash_hasher};
4356 std::unordered_set<uint256, SaltedUint256Hasher> seen_wtxids{0, m_txhash_hasher};
4360 const auto current_time{GetTime<std::chrono::microseconds>()};
4363 for (
CInv& inv : vInv) {
4364 if (interruptMsgProc)
return;
4369 if (peer.m_wtxid_relay) {
4376 const bool fAlreadyHave = AlreadyHaveBlock(inv.
hash);
4379 UpdateBlockAvailability(pfrom.
GetId(), inv.
hash);
4387 best_block = &inv.
hash;
4390 if (reject_tx_invs) {
4396 auto& seen_hashes{inv.
IsMsgWtx() ? seen_wtxids : seen_txids};
4397 if (!seen_hashes.insert(inv.
hash).second)
continue;
4399 AddKnownTx(peer, inv.
hash);
4402 const bool fAlreadyHave{m_txdownloadman.AddTxAnnouncement(pfrom.
GetId(), gtxid, current_time)};
4410 if (best_block !=
nullptr) {
4422 if (state.fSyncStarted || (!peer.m_inv_triggered_getheaders_before_sync && *best_block != m_last_block_inv_triggering_headers_sync)) {
4423 if (MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer)) {
4425 m_chainman.m_best_header->nHeight, best_block->ToString(),
4428 if (!state.fSyncStarted) {
4429 peer.m_inv_triggered_getheaders_before_sync =
true;
4433 m_last_block_inv_triggering_headers_sync = *best_block;
4442 std::vector<CInv> vInv;
4446 Misbehaving(peer,
strprintf(
"getdata message size = %u", vInv.size()));
4452 if (vInv.size() > 0) {
4457 const auto pushed_tx_opt{m_tx_for_private_broadcast.GetTxForNode(pfrom.
GetId())};
4458 if (!pushed_tx_opt) {
4469 if (vInv.size() == 1 && vInv[0].IsMsgTx() && vInv[0].hash == pushed_tx->GetHash().ToUint256()) {
4473 peer.m_ping_queued =
true;
4484 LOCK(peer.m_getdata_requests_mutex);
4485 peer.m_getdata_requests.insert(peer.m_getdata_requests.end(), vInv.begin(), vInv.end());
4486 ProcessGetData(pfrom, peer, interruptMsgProc);
4495 vRecv >> locator >> hashStop;
4511 std::shared_ptr<const CBlock> a_recent_block;
4513 LOCK(m_most_recent_block_mutex);
4514 a_recent_block = m_most_recent_block;
4517 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
4532 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4547 if (--nLimit <= 0) {
4551 WITH_LOCK(peer.m_block_inv_mutex, {peer.m_continuation_block = pindex->GetBlockHash();});
4571 for (
size_t i = 1; i < req.
indexes.size(); ++i) {
4575 std::shared_ptr<const CBlock> recent_block;
4577 LOCK(m_most_recent_block_mutex);
4578 if (m_most_recent_block_hash == req.
blockhash)
4579 recent_block = m_most_recent_block;
4583 SendBlockTransactions(pfrom, peer, *recent_block, req);
4602 if (!block_pos.IsNull()) {
4609 SendBlockTransactions(pfrom, peer, block, req);
4622 WITH_LOCK(peer.m_getdata_requests_mutex, peer.m_getdata_requests.push_back(inv));
4630 vRecv >> locator >> hashStop;
4648 if (m_chainman.
ActiveTip() ==
nullptr ||
4650 LogDebug(
BCLog::NET,
"Ignoring getheaders from peer=%d because active chain has too little work; sending empty response\n", pfrom.
GetId());
4657 CNodeState *nodestate =
State(pfrom.
GetId());
4666 if (!BlockRequestAllowed(*pindex)) {
4667 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block header that isn't in the main chain\n", __func__, pfrom.
GetId());
4680 std::vector<CBlock> vHeaders;
4681 int nLimit = m_opts.max_headers_result;
4683 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4686 if (--nLimit <= 0 || pindex->GetBlockHash() == hashStop)
4701 nodestate->pindexBestHeaderSent = pindex ? pindex : m_chainman.
ActiveChain().
Tip();
4707 if (RejectIncomingTxs(pfrom)) {
4721 const Txid& txid = ptx->GetHash();
4722 const Wtxid& wtxid = ptx->GetWitnessHash();
4725 AddKnownTx(peer, hash);
4727 if (
const auto num_broadcasted{m_tx_for_private_broadcast.Remove(ptx)}) {
4729 "network from %s; stopping private broadcast attempts",
4740 const auto& [should_validate, package_to_validate] = m_txdownloadman.ReceivedTx(pfrom.
GetId(), ptx);
4741 if (!should_validate) {
4746 if (!m_mempool.
exists(txid)) {
4747 LogInfo(
"Not relaying non-mempool transaction %s (wtxid=%s) from forcerelay peer=%d\n",
4750 LogInfo(
"Force relaying tx %s (wtxid=%s) from peer=%d\n",
4752 InitiateTxBroadcastToAll(wtxid);
4756 if (package_to_validate) {
4759 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4760 ProcessPackageResult(package_to_validate.value(), package_result);
4766 Assume(!package_to_validate.has_value());
4776 if (
auto package_to_validate{ProcessInvalidTx(pfrom.
GetId(), ptx, state,
true)}) {
4779 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4780 ProcessPackageResult(package_to_validate.value(), package_result);
4793 }
else if (m_opts.ignore_incoming_txs) {
4800 const CNodeState *nodestate =
State(pfrom.
GetId());
4801 if (!nodestate->m_provides_cmpctblocks) {
4808 vRecv >> cmpctblock;
4810 bool received_new_header =
false;
4820 MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer);
4830 received_new_header =
true;
4838 MaybePunishNodeForBlock(pfrom.
GetId(), state,
true,
"invalid header via cmpctblock");
4845 if (received_new_header) {
4846 LogBlockHeader(*pindex, pfrom,
true);
4849 bool fProcessBLOCKTXN =
false;
4853 bool fRevertToHeaderProcessing =
false;
4857 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
4858 bool fBlockReconstructed =
false;
4864 CNodeState *nodestate =
State(pfrom.
GetId());
4875 auto range_flight = mapBlocksInFlight.equal_range(pindex->
GetBlockHash());
4876 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
4877 bool requested_block_from_this_peer{
false};
4880 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
4882 while (range_flight.first != range_flight.second) {
4883 if (range_flight.first->second.first == pfrom.
GetId()) {
4884 requested_block_from_this_peer =
true;
4887 range_flight.first++;
4897 if (requested_block_from_this_peer) {
4900 std::vector<CInv> vInv(1);
4908 if (!already_in_flight && !CanDirectFetch()) {
4916 requested_block_from_this_peer) {
4917 std::list<QueuedBlock>::iterator* queuedBlockIt =
nullptr;
4918 if (!BlockRequested(pfrom.
GetId(), *pindex, &queuedBlockIt)) {
4919 if (!(*queuedBlockIt)->partialBlock)
4932 Misbehaving(peer,
"invalid compact block");
4935 if (first_in_flight) {
4937 std::vector<CInv> vInv(1);
4948 for (
size_t i = 0; i < cmpctblock.
BlockTxCount(); i++) {
4953 fProcessBLOCKTXN =
true;
4954 }
else if (first_in_flight) {
4961 IsBlockRequestedFromOutbound(blockhash) ||
4980 ReadStatus status = tempBlock.InitData(cmpctblock, vExtraTxnForCompact);
4985 std::vector<CTransactionRef> dummy;
4987 status = tempBlock.FillBlock(*pblock, dummy,
4990 fBlockReconstructed =
true;
4994 if (requested_block_from_this_peer) {
4997 std::vector<CInv> vInv(1);
5003 fRevertToHeaderProcessing =
true;
5008 if (fProcessBLOCKTXN) {
5011 return ProcessCompactBlockTxns(pfrom, peer, txn);
5014 if (fRevertToHeaderProcessing) {
5020 return ProcessHeadersMessage(pfrom, peer, {cmpctblock.
header},
true);
5023 if (fBlockReconstructed) {
5028 mapBlockSource.emplace(pblock->GetHash(), std::make_pair(pfrom.
GetId(),
false));
5046 RemoveBlockRequest(pblock->GetHash(), std::nullopt);
5063 return ProcessCompactBlockTxns(pfrom, peer, resp);
5074 std::vector<CBlockHeader>
headers;
5078 if (nCount > m_opts.max_headers_result) {
5079 Misbehaving(peer,
strprintf(
"headers message size = %u", nCount));
5083 for (
unsigned int n = 0; n < nCount; n++) {
5088 ProcessHeadersMessage(pfrom, peer, std::move(
headers),
false);
5092 if (m_headers_presync_should_signal.exchange(
false)) {
5093 HeadersPresyncStats stats;
5095 LOCK(m_headers_presync_mutex);
5096 auto it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
5097 if (it != m_headers_presync_stats.end()) stats = it->second;
5115 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
5126 Misbehaving(peer,
"mutated block");
5131 bool forceProcessing =
false;
5132 const uint256 hash(pblock->GetHash());
5133 bool min_pow_checked =
false;
5138 forceProcessing = IsBlockRequested(hash);
5139 RemoveBlockRequest(hash, pfrom.
GetId());
5143 mapBlockSource.emplace(hash, std::make_pair(pfrom.
GetId(),
true));
5147 min_pow_checked =
true;
5150 ProcessBlock(pfrom, pblock, forceProcessing, min_pow_checked);
5167 Assume(SetupAddressRelay(pfrom, peer));
5171 if (peer.m_getaddr_recvd) {
5175 peer.m_getaddr_recvd =
true;
5177 peer.m_addrs_to_send.clear();
5178 std::vector<CAddress> vAddr;
5184 for (
const CAddress &addr : vAddr) {
5185 PushAddress(peer, addr);
5213 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5214 LOCK(tx_relay->m_tx_inventory_mutex);
5215 tx_relay->m_send_mempool =
true;
5241 ProcessPong(pfrom, peer, time_received, vRecv);
5257 Misbehaving(peer,
"too-large bloom filter");
5258 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5260 LOCK(tx_relay->m_bloom_filter_mutex);
5261 tx_relay->m_bloom_filter.reset(
new CBloomFilter(filter));
5262 tx_relay->m_relay_txs =
true;
5266 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5277 std::vector<unsigned char> vData;
5285 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5286 LOCK(tx_relay->m_bloom_filter_mutex);
5287 if (tx_relay->m_bloom_filter) {
5288 tx_relay->m_bloom_filter->insert(vData);
5294 Misbehaving(peer,
"bad filteradd message");
5305 auto tx_relay = peer.GetTxRelay();
5306 if (!tx_relay)
return;
5309 LOCK(tx_relay->m_bloom_filter_mutex);
5310 tx_relay->m_bloom_filter =
nullptr;
5311 tx_relay->m_relay_txs =
true;
5315 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5321 vRecv >> newFeeFilter;
5323 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5324 tx_relay->m_fee_filter_received = newFeeFilter;
5332 ProcessGetCFilters(pfrom, peer, vRecv);
5337 ProcessGetCFHeaders(pfrom, peer, vRecv);
5342 ProcessGetCFCheckPt(pfrom, peer, vRecv);
5347 std::vector<CInv> vInv;
5349 std::vector<GenTxid> tx_invs;
5351 for (
CInv &inv : vInv) {
5357 LOCK(m_tx_download_mutex);
5358 m_txdownloadman.ReceivedNotFound(pfrom.
GetId(), tx_invs);
5367bool PeerManagerImpl::MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer)
5370 LOCK(peer.m_misbehavior_mutex);
5373 if (!peer.m_should_discourage)
return false;
5375 peer.m_should_discourage =
false;
5380 LogWarning(
"Not punishing noban peer %d!", peer.m_id);
5386 LogWarning(
"Not punishing manually connected peer %d!", peer.m_id);
5412bool PeerManagerImpl::MaybeDisconnectForTxRelayCapacity(
CNode&
node,
const std::string& msg_type, std::optional<NodeId> protect_peer)
5414 if (!
node.IsInboundConn() || !
node.m_relays_txs)
return false;
5417 LogDebug(
BCLog::NET,
"failed to find a tx-relaying eviction candidate - connection dropped after %s message, peer=%d\n", msg_type,
node.GetId());
5418 node.fDisconnect =
true;
5422bool PeerManagerImpl::ProcessMessages(
CNode&
node, std::atomic<bool>& interruptMsgProc)
5427 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
5428 if (maybe_peer ==
nullptr)
return false;
5429 Peer& peer{*maybe_peer};
5433 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent)
return false;
5436 LOCK(peer.m_getdata_requests_mutex);
5437 if (!peer.m_getdata_requests.empty()) {
5438 ProcessGetData(
node, peer, interruptMsgProc);
5442 const bool processed_orphan = ProcessOrphanTx(peer);
5444 if (
node.fDisconnect)
5447 if (processed_orphan)
return true;
5452 LOCK(peer.m_getdata_requests_mutex);
5453 if (!peer.m_getdata_requests.empty())
return true;
5457 if (
node.fPauseSend)
return false;
5459 auto poll_result{
node.PollMessage()};
5466 bool fMoreWork = poll_result->second;
5470 node.m_addr_name.c_str(),
5471 node.ConnectionTypeAsString().c_str(),
5477 if (m_opts.capture_messages) {
5482 ProcessMessage(peer,
node,
msg.m_type,
msg.m_recv,
msg.m_time, interruptMsgProc);
5483 if (interruptMsgProc)
return false;
5485 LOCK(peer.m_getdata_requests_mutex);
5486 if (!peer.m_getdata_requests.empty()) fMoreWork =
true;
5493 LOCK(m_tx_download_mutex);
5494 if (m_txdownloadman.HaveMoreWork(peer.m_id)) fMoreWork =
true;
5495 }
catch (
const std::exception& e) {
5504void PeerManagerImpl::ConsiderEviction(
CNode& pto, Peer& peer, std::chrono::seconds time_in_seconds)
5517 if (state.pindexBestKnownBlock !=
nullptr && state.pindexBestKnownBlock->nChainWork >= m_chainman.
ActiveChain().
Tip()->
nChainWork) {
5519 if (state.m_chain_sync.m_timeout != 0
s) {
5520 state.m_chain_sync.m_timeout = 0
s;
5521 state.m_chain_sync.m_work_header =
nullptr;
5522 state.m_chain_sync.m_sent_getheaders =
false;
5524 }
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)) {
5532 state.m_chain_sync.m_work_header = m_chainman.
ActiveChain().
Tip();
5533 state.m_chain_sync.m_sent_getheaders =
false;
5534 }
else if (state.m_chain_sync.m_timeout > 0
s && time_in_seconds > state.m_chain_sync.m_timeout) {
5538 if (state.m_chain_sync.m_sent_getheaders) {
5540 LogInfo(
"Outbound peer has old chain, best known block = %s, %s", state.pindexBestKnownBlock !=
nullptr ? state.pindexBestKnownBlock->GetBlockHash().ToString() :
"<none>", pto.
DisconnectMsg());
5543 assert(state.m_chain_sync.m_work_header);
5548 MaybeSendGetHeaders(pto,
5549 GetLocator(state.m_chain_sync.m_work_header->pprev),
5551 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());
5552 state.m_chain_sync.m_sent_getheaders =
true;
5573 std::pair<NodeId, std::chrono::seconds> youngest_peer{-1, 0}, next_youngest_peer{-1, 0};
5577 if (pnode->
GetId() > youngest_peer.first) {
5578 next_youngest_peer = youngest_peer;
5579 youngest_peer.first = pnode->GetId();
5580 youngest_peer.second = pnode->m_last_block_time;
5583 NodeId to_disconnect = youngest_peer.first;
5584 if (youngest_peer.second > next_youngest_peer.second) {
5587 to_disconnect = next_youngest_peer.first;
5596 CNodeState *node_state =
State(pnode->
GetId());
5597 if (node_state ==
nullptr ||
5600 LogDebug(
BCLog::NET,
"disconnecting extra block-relay-only peer=%d (last block received at time %d)\n",
5604 LogDebug(
BCLog::NET,
"keeping block-relay-only peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5605 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), node_state->vBlocksInFlight.size());
5623 std::optional<WorstPeer> worst_peer;
5632 if (state ==
nullptr)
return;
5634 if (state->m_chain_sync.m_protect)
return;
5638 if (!worst_peer.has_value() ||
5639 (state->m_last_block_announcement < (*worst_peer).oldest_block_announcement) ||
5640 ((state->m_last_block_announcement == (*worst_peer).oldest_block_announcement) && pnode->
GetId() > (*worst_peer).node)) {
5641 worst_peer = WorstPeer{pnode->
GetId(), state->m_last_block_announcement};
5644 if (worst_peer.has_value()) {
5655 LogDebug(
BCLog::NET,
"disconnecting extra outbound peer=%d (last block announcement received at time %d)\n",
5656 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>((*worst_peer).oldest_block_announcement));
5660 LogDebug(
BCLog::NET,
"keeping outbound peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5661 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), state.vBlocksInFlight.size());
5677void PeerManagerImpl::CheckForStaleTipAndEvictPeers()
5682 auto now{GetTime<std::chrono::seconds>()};
5684 EvictExtraOutboundPeers(current_time);
5686 if (now > m_stale_tip_check_time) {
5690 LogInfo(
"Potential stale tip detected, will try using extra outbound peer (last tip update: %d seconds ago)\n",
5699 if (!m_initial_sync_finished && CanDirectFetch()) {
5701 m_initial_sync_finished =
true;
5708 peer.m_ping_nonce_sent &&
5718 bool pingSend =
false;
5720 if (peer.m_ping_queued) {
5725 if (peer.m_ping_nonce_sent == 0 && now > peer.m_ping_start.load() +
PING_INTERVAL) {
5734 }
while (
nonce == 0);
5735 peer.m_ping_queued =
false;
5736 peer.m_ping_start = now;
5738 peer.m_ping_nonce_sent =
nonce;
5742 peer.m_ping_nonce_sent = 0;
5748void PeerManagerImpl::MaybeSendAddr(
CNode&
node, Peer& peer, std::chrono::microseconds current_time)
5751 if (!peer.m_addr_relay_enabled)
return;
5753 LOCK(peer.m_addr_send_times_mutex);
5756 peer.m_next_local_addr_send < current_time) {
5763 if (peer.m_next_local_addr_send != 0us) {
5764 peer.m_addr_known->reset();
5767 CAddress local_addr{*local_service, peer.m_our_services, Now<NodeSeconds>()};
5768 if (peer.m_next_local_addr_send == 0us) {
5772 if (IsAddrCompatible(peer, local_addr)) {
5773 std::vector<CAddress> self_announcement{local_addr};
5774 if (peer.m_wants_addrv2) {
5782 PushAddress(peer, local_addr);
5789 if (current_time <= peer.m_next_addr_send)
return;
5802 bool ret = peer.m_addr_known->contains(addr.
GetKey());
5803 if (!
ret) peer.m_addr_known->insert(addr.
GetKey());
5806 peer.m_addrs_to_send.erase(std::remove_if(peer.m_addrs_to_send.begin(), peer.m_addrs_to_send.end(), addr_already_known),
5807 peer.m_addrs_to_send.end());
5810 if (peer.m_addrs_to_send.empty())
return;
5812 if (peer.m_wants_addrv2) {
5817 peer.m_addrs_to_send.clear();
5820 if (peer.m_addrs_to_send.capacity() > 40) {
5821 peer.m_addrs_to_send.shrink_to_fit();
5825void PeerManagerImpl::MaybeSendSendHeaders(
CNode&
node, Peer& peer)
5833 CNodeState &state = *
State(
node.GetId());
5834 if (state.pindexBestKnownBlock !=
nullptr &&
5841 peer.m_sent_sendheaders =
true;
5846void PeerManagerImpl::MaybeSendFeefilter(
CNode& pto, Peer& peer, std::chrono::microseconds current_time)
5848 if (m_opts.ignore_incoming_txs)
return;
5864 if (peer.m_fee_filter_sent == MAX_FILTER) {
5867 peer.m_next_send_feefilter = 0us;
5870 if (current_time > peer.m_next_send_feefilter) {
5871 CAmount filterToSend = m_fee_filter_rounder.round(currentFilter);
5874 if (filterToSend != peer.m_fee_filter_sent) {
5876 peer.m_fee_filter_sent = filterToSend;
5883 (currentFilter < 3 * peer.m_fee_filter_sent / 4 || currentFilter > 4 * peer.m_fee_filter_sent / 3)) {
5888bool PeerManagerImpl::RejectIncomingTxs(
const CNode& peer)
const
5901 const size_t nAvail{vRecv.
size()};
5902 bool bPingFinished =
false;
5903 std::string sProblem;
5905 if (nAvail >=
sizeof(
nonce)) {
5909 if (peer.m_ping_nonce_sent != 0) {
5910 if (
nonce == peer.m_ping_nonce_sent) {
5912 bPingFinished =
true;
5913 const auto ping_time = ping_end - peer.m_ping_start.load();
5914 if (ping_time.count() >= 0) {
5918 m_tx_for_private_broadcast.NodeConfirmedReception(pfrom.
GetId());
5925 sProblem =
"Timing mishap";
5929 sProblem =
"Nonce mismatch";
5932 bPingFinished =
true;
5933 sProblem =
"Nonce zero";
5937 sProblem =
"Unsolicited pong without ping";
5941 bPingFinished =
true;
5942 sProblem =
"Short payload";
5945 if (!(sProblem.empty())) {
5949 peer.m_ping_nonce_sent,
5953 if (bPingFinished) {
5954 peer.m_ping_nonce_sent = 0;
5958bool PeerManagerImpl::SetupAddressRelay(
const CNode&
node, Peer& peer)
5963 if (
node.IsBlockOnlyConn())
return false;
5968 if (
node.IsFeelerConn())
return false;
5970 if (!peer.m_addr_relay_enabled.exchange(
true)) {
5974 peer.m_addr_known = std::make_unique<CRollingBloomFilter>(5000, 0.001);
5980void PeerManagerImpl::ProcessAddrs(std::string_view msg_type,
CNode& pfrom, Peer& peer, std::vector<CAddress>&& vAddr,
const std::atomic<bool>& interruptMsgProc)
5985 if (!SetupAddressRelay(pfrom, peer)) {
5992 Misbehaving(peer,
strprintf(
"%s message size = %u", msg_type, vAddr.size()));
5997 std::vector<CAddress> vAddrOk;
6003 const auto time_diff{current_time - peer.m_addr_token_timestamp};
6007 peer.m_addr_token_timestamp = current_time;
6010 uint64_t num_proc = 0;
6011 uint64_t num_rate_limit = 0;
6012 std::shuffle(vAddr.begin(), vAddr.end(), m_rng);
6015 if (interruptMsgProc)
6019 if (peer.m_addr_token_bucket < 1.0) {
6025 peer.m_addr_token_bucket -= 1.0;
6034 addr.
nTime = std::chrono::time_point_cast<std::chrono::seconds>(current_time - 5 * 24h);
6036 AddAddressKnown(peer, addr);
6043 if (addr.
nTime > current_time - 10min && vAddr.size() <= 10 && addr.
IsRoutable()) {
6045 RelayAddress(pfrom.
GetId(), addr, reachable);
6049 vAddrOk.push_back(addr);
6052 peer.m_addr_processed += num_proc;
6053 peer.m_addr_rate_limited += num_rate_limit;
6054 LogDebug(
BCLog::NET,
"Received addr: %u addresses (%u processed, %u rate-limited) from peer=%d\n",
6055 vAddr.size(), num_proc, num_rate_limit, pfrom.
GetId());
6057 m_addrman.
Add(vAddrOk, pfrom.
addr, 2h);
6066bool PeerManagerImpl::SendMessages(
CNode&
node)
6071 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
6072 if (!maybe_peer)
return false;
6073 Peer& peer{*maybe_peer};
6078 if (MaybeDiscourageAndDisconnect(
node, peer))
return true;
6081 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent) {
6082 PushNodeVersion(
node, peer);
6083 peer.m_outbound_version_message_sent =
true;
6087 if (!
node.fSuccessfullyConnected ||
node.fDisconnect)
6091 const auto current_time{GetTime<std::chrono::microseconds>()};
6096 if (
node.IsPrivateBroadcastConn()) {
6100 node.fDisconnect =
true;
6107 node.fDisconnect =
true;
6111 MaybeSendPing(
node, peer, now);
6114 if (
node.fDisconnect)
return true;
6116 MaybeSendAddr(
node, peer, current_time);
6118 MaybeSendSendHeaders(
node, peer);
6120 ProcessInvBacklog(now);
6125 CNodeState &state = *
State(
node.GetId());
6128 if (m_chainman.m_best_header ==
nullptr) {
6135 bool sync_blocks_and_headers_from_peer =
false;
6136 if (state.fPreferredDownload) {
6137 sync_blocks_and_headers_from_peer =
true;
6138 }
else if (CanServeBlocks(peer) && !
node.IsAddrFetchConn()) {
6148 if (m_num_preferred_download_peers == 0 || mapBlocksInFlight.empty()) {
6149 sync_blocks_and_headers_from_peer =
true;
6155 if ((nSyncStarted == 0 && sync_blocks_and_headers_from_peer) || m_chainman.m_best_header->Time() >
NodeClock::now() - 24h) {
6156 const CBlockIndex* pindexStart = m_chainman.m_best_header;
6164 if (pindexStart->
pprev)
6165 pindexStart = pindexStart->
pprev;
6169 state.fSyncStarted =
true;
6193 LOCK(peer.m_block_inv_mutex);
6194 std::vector<CBlock> vHeaders;
6195 bool fRevertToInv = ((!peer.m_prefers_headers &&
6196 (!state.m_requested_hb_cmpctblocks || peer.m_blocks_for_headers_relay.size() > 1)) ||
6199 ProcessBlockAvailability(
node.GetId());
6201 if (!fRevertToInv) {
6202 bool fFoundStartingHeader =
false;
6206 for (
const uint256& hash : peer.m_blocks_for_headers_relay) {
6211 fRevertToInv =
true;
6214 if (pBestIndex !=
nullptr && pindex->
pprev != pBestIndex) {
6226 fRevertToInv =
true;
6229 pBestIndex = pindex;
6230 if (fFoundStartingHeader) {
6233 }
else if (PeerHasHeader(&state, pindex)) {
6235 }
else if (pindex->
pprev ==
nullptr || PeerHasHeader(&state, pindex->
pprev)) {
6238 fFoundStartingHeader =
true;
6243 fRevertToInv =
true;
6248 if (!fRevertToInv && !vHeaders.empty()) {
6249 if (vHeaders.size() == 1 && state.m_requested_hb_cmpctblocks) {
6253 vHeaders.front().GetHash().ToString(),
node.GetId());
6255 std::optional<CSerializedNetMsg> cached_cmpctblock_msg;
6257 LOCK(m_most_recent_block_mutex);
6258 if (m_most_recent_block_hash == pBestIndex->
GetBlockHash()) {
6262 if (cached_cmpctblock_msg.has_value()) {
6263 PushMessage(
node, std::move(cached_cmpctblock_msg.value()));
6271 state.pindexBestHeaderSent = pBestIndex;
6272 }
else if (peer.m_prefers_headers) {
6273 if (vHeaders.size() > 1) {
6276 vHeaders.front().GetHash().ToString(),
6277 vHeaders.back().GetHash().ToString(),
node.GetId());
6280 vHeaders.front().GetHash().ToString(),
node.GetId());
6283 state.pindexBestHeaderSent = pBestIndex;
6285 fRevertToInv =
true;
6291 if (!peer.m_blocks_for_headers_relay.empty()) {
6292 const uint256& hashToAnnounce = peer.m_blocks_for_headers_relay.back();
6305 if (!PeerHasHeader(&state, pindex)) {
6306 peer.m_blocks_for_inv_relay.push_back(hashToAnnounce);
6312 peer.m_blocks_for_headers_relay.clear();
6318 std::vector<CInv> vInv;
6320 LOCK(peer.m_block_inv_mutex);
6321 vInv.reserve(peer.m_blocks_for_inv_relay.size());
6324 for (
const uint256& hash : peer.m_blocks_for_inv_relay) {
6331 peer.m_blocks_for_inv_relay.clear();
6334 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
6335 LOCK(tx_relay->m_tx_inventory_mutex);
6338 if (tx_relay->m_next_inv_send_time < current_time) {
6339 fSendTrickle =
true;
6340 if (
node.IsInboundConn()) {
6349 LOCK(tx_relay->m_bloom_filter_mutex);
6350 if (!tx_relay->m_relay_txs) tx_relay->m_tx_inventory_to_send.clear();
6354 if (fSendTrickle && tx_relay->m_send_mempool) {
6355 auto vtxinfo = m_mempool.
infoAll();
6360 tx_relay->m_send_mempool =
false;
6361 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6364 tx_relay->m_tx_inventory_to_send.clear();
6366 LOCK(tx_relay->m_bloom_filter_mutex);
6368 for (
const auto& txinfo : vtxinfo) {
6369 const Txid& txid{txinfo.tx->GetHash()};
6370 const Wtxid& wtxid{txinfo.tx->GetWitnessHash()};
6371 const auto inv = peer.m_wtxid_relay ?
6376 if (txinfo.fee < filterrate.GetFee(txinfo.vsize)) {
6379 if (tx_relay->m_bloom_filter) {
6380 if (!tx_relay->m_bloom_filter->IsRelevantAndUpdate(*txinfo.tx))
continue;
6382 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6383 vInv.push_back(inv);
6395 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6398 auto& invs = tx_relay->m_tx_inventory_to_send;
6399 std::vector<CTransactionRef> res;
6401 if (invs.size() == 0)
return res;
6404 if (invs.capacity() > 2 * invs.size()) invs.shrink_to_fit();
6408 res.reserve(txiters.size());
6409 for (
auto txiter : txiters) {
6410 if (txiter->GetFee() < filterrate.GetFee(txiter->GetTxSize())) {
6413 res.push_back(txiter->GetSharedTx());
6416 tx_relay->m_last_inv_sequence = m_mempool.
GetSequence();
6420 LOCK(tx_relay->m_bloom_filter_mutex);
6421 vInv.reserve(std::min<size_t>(
MAX_INV_SZ, vInv.size() + inv_tx.size()));
6422 for (
auto& tx : inv_tx) {
6426 const auto inv = peer.m_wtxid_relay ?
6430 if (tx_relay->m_tx_inventory_known_filter.contains(inv.
hash)) {
6433 if (tx_relay->m_bloom_filter && !tx_relay->m_bloom_filter->IsRelevantAndUpdate(*tx))
continue;
6435 vInv.push_back(inv);
6440 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6448 auto stalling_timeout = m_block_stalling_timeout.load();
6449 if (state.m_stalling_since.count() && state.m_stalling_since < current_time - stalling_timeout) {
6453 if (
node.IsManualConn()) {
6456 while (!state.vBlocksInFlight.empty()) {
6457 RemoveBlockRequest(state.vBlocksInFlight.front().pindex->GetBlockHash(),
node.GetId());
6460 LogInfo(
"Peer is stalling block download, %s",
node.DisconnectMsg());
6461 node.fDisconnect =
true;
6466 if (stalling_timeout != new_timeout && m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
6476 if (state.vBlocksInFlight.size() > 0) {
6477 QueuedBlock &queuedBlock = state.vBlocksInFlight.front();
6478 int nOtherPeersWithValidatedDownloads = m_peers_downloading_from - 1;
6480 LogInfo(
"Timeout downloading block %s, %s", queuedBlock.pindex->GetBlockHash().ToString(),
node.DisconnectMsg());
6481 node.fDisconnect =
true;
6486 if (state.fSyncStarted && peer.m_headers_sync_timeout < std::chrono::microseconds::max()) {
6488 if (m_chainman.m_best_header->Time() <=
NodeClock::now() - 24h) {
6489 if (current_time > peer.m_headers_sync_timeout && nSyncStarted == 1 && (m_num_preferred_download_peers - state.fPreferredDownload >= 1)) {
6496 LogInfo(
"Timeout downloading headers, %s",
node.DisconnectMsg());
6497 node.fDisconnect =
true;
6500 LogInfo(
"Timeout downloading headers from noban peer, not %s",
node.DisconnectMsg());
6506 state.fSyncStarted =
false;
6508 peer.m_headers_sync_timeout = 0us;
6514 peer.m_headers_sync_timeout = std::chrono::microseconds::max();
6520 ConsiderEviction(
node, peer, GetTime<std::chrono::seconds>());
6525 std::vector<CInv> vGetData;
6526 const bool can_request_blocks_from_peer{current_time >= state.m_block_download_paused_until};
6528 std::vector<const CBlockIndex*> vToDownload;
6530 auto get_inflight_budget = [&state]() {
6537 FindNextBlocksToDownload(peer, get_inflight_budget(), vToDownload, staller);
6538 auto historical_blocks{m_chainman.GetHistoricalBlockRange()};
6539 if (historical_blocks && !IsLimitedPeer(peer)) {
6543 TryDownloadingHistoricalBlocks(
6545 get_inflight_budget(),
6546 vToDownload, from_tip, historical_blocks->second);
6549 uint32_t nFetchFlags = GetFetchFlags(peer);
6551 BlockRequested(
node.GetId(), *pindex);
6555 if (state.vBlocksInFlight.empty() && staller != -1) {
6556 if (
State(staller)->m_stalling_since == 0us) {
6557 State(staller)->m_stalling_since = current_time;
6567 LOCK(m_tx_download_mutex);
6568 for (
const GenTxid& gtxid : m_txdownloadman.GetRequestsToSend(
node.GetId(), current_time)) {
6577 if (!vGetData.empty())
6580 MaybeSendFeefilter(
node, peer, current_time);
constexpr CAmount MAX_MONEY
No amount larger than this (in satoshi) is valid.
bool MoneyRange(const CAmount &nValue)
int64_t CAmount
Amount in satoshis (Can be negative)
enum ReadStatus_t ReadStatus
const std::string & BlockFilterTypeName(BlockFilterType filter_type)
Get the human-readable name for a filter type.
BlockFilterIndex * GetBlockFilterIndex(BlockFilterType filter_type)
Get a block filter index by type.
constexpr int CFCHECKPT_INTERVAL
Interval between compact filter checkpoints.
CBlockLocator GetLocator(const CBlockIndex *index)
Get a locator for a block index entry.
int64_t GetBlockProofEquivalentTime(const CBlockIndex &to, const CBlockIndex &from, const CBlockIndex &tip, const Consensus::Params ¶ms)
Return the time it would take to redo the work difference between from and to, assuming the current h...
const CBlockIndex * LastCommonAncestor(const CBlockIndex *pa, const CBlockIndex *pb)
Find the last common ancestor two blocks have.
@ BLOCK_VALID_CHAIN
Outputs do not overspend inputs, no double spends, coinbase output ok, no immature coinbase spends,...
@ BLOCK_VALID_TRANSACTIONS
Only first tx is coinbase, 2 <= coinbase input script length <= 100, transactions valid,...
@ BLOCK_VALID_SCRIPTS
Scripts & signatures ok.
@ BLOCK_VALID_TREE
All parent headers found, difficulty matches, timestamp >= median previous.
@ BLOCK_HAVE_DATA
full block available in blk*.dat
arith_uint256 GetBlockProof(const CBlockIndex &block)
Compute how much work a block index entry corresponds to.
#define Assert(val)
Identity function.
#define Assume(val)
Assume is the identity function.
Stochastic address manager.
void Connected(const CService &addr, NodeSeconds time=Now< NodeSeconds >())
We have successfully connected to this peer.
bool Good(const CService &addr, NodeSeconds time=Now< NodeSeconds >())
Mark an address record as accessible and attempt to move it to addrman's tried table.
bool Add(const std::vector< CAddress > &vAddr, const CNetAddr &source, std::chrono::seconds time_penalty=0s)
Attempt to add one or more addresses to addrman's new table.
void SetServices(const CService &addr, ServiceFlags nServices)
Update an entry's service bits.
bool IsBanned(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
Return whether net_addr is banned.
bool IsDiscouraged(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
Return whether net_addr is discouraged.
void Discourage(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
BlockFilterIndex is used to store and retrieve block filters, hashes, and headers for a range of bloc...
bool LookupFilterRange(int start_height, const CBlockIndex *stop_index, std::vector< BlockFilter > &filters_out) const
Get a range of filters between two heights on a chain.
bool LookupFilterHashRange(int start_height, const CBlockIndex *stop_index, std::vector< uint256 > &hashes_out) const
Get a range of filter hashes between two heights on a chain.
bool LookupFilterHeader(const CBlockIndex *block_index, uint256 &header_out) EXCLUSIVE_LOCKS_REQUIRED(!m_cs_headers_cache)
Get a single filter header by block.
std::vector< CTransactionRef > txn
std::vector< uint16_t > indexes
A CService with information about it as peer.
ServiceFlags nServices
Serialized as uint64_t in V1, and as CompactSize in V2.
static constexpr SerParams V1_NETWORK
NodeSeconds nTime
Always included in serialization. The behavior is unspecified if the value is not representable as ui...
static constexpr SerParams V2_NETWORK
size_t BlockTxCount() const
std::vector< CTransactionRef > vtx
The block chain is a tree shaped structure starting with the genesis block at the root,...
bool IsValid(enum BlockStatus nUpTo) const EXCLUSIVE_LOCKS_REQUIRED(
Check whether this block index entry is valid up to the passed validity level.
CBlockIndex * pprev
pointer to the index of the predecessor of this block
CBlockHeader GetBlockHeader() const
arith_uint256 nChainWork
(memory only) Total amount of work (expected number of hashes) in the chain up to and including this ...
bool HaveNumChainTxs() const
Check whether this block and all previous blocks back to the genesis block or an assumeutxo snapshot ...
uint256 GetBlockHash() const
int64_t GetBlockTime() const
unsigned int nTx
Number of transactions in this block.
CBlockIndex * GetAncestor(int height)
Efficiently find an ancestor of this block.
int nHeight
height of the entry in the chain. The genesis block has height 0
FlatFilePos GetBlockPos() const EXCLUSIVE_LOCKS_REQUIRED(
BloomFilter is a probabilistic filter which SPV clients provide so that we can filter the transaction...
bool IsWithinSizeConstraints() const
True if the size is <= MAX_BLOOM_FILTER_SIZE and the number of hash functions is <= MAX_HASH_FUNCS (c...
An in-memory indexed chain of blocks.
bool Contains(const CBlockIndex &index) const
Efficiently check whether a block is present in this chain.
CBlockIndex * Tip() const
Returns the index entry for the tip of this chain, or nullptr if none.
CBlockIndex * Next(const CBlockIndex &index) const
Find the successor of a block in this chain, or nullptr if the given index is not found or is the tip...
int Height() const
Return the maximal height in the chain.
CChainParams defines various tweakable parameters of a given instance of the Bitcoin system.
const HeadersSyncParams & HeadersSync() const
const Consensus::Params & GetConsensus() const
void NumToOpenAdd(size_t n)
Increment the number of new connections of type ConnectionType::PRIVATE_BROADCAST to be opened by CCo...
size_t NumToOpenSub(size_t n)
Decrement the number of new connections of type ConnectionType::PRIVATE_BROADCAST to be opened by CCo...
bool GetNetworkActive() const
bool GetTryNewOutboundPeer() const
class CConnman::PrivateBroadcast m_private_broadcast
bool ShouldRunInactivityChecks(const CNode &node, NodeClock::time_point now) const
Return true if we should disconnect the peer for failing an inactivity check.
std::vector< CAddress > GetAddresses(CNode &requestor, size_t max_addresses, size_t max_pct)
Return addresses from the per-requestor cache.
void SetTryNewOutboundPeer(bool flag)
void WakeMessageHandler() EXCLUSIVE_LOCKS_REQUIRED(!mutexMsgProc)
bool OutboundTargetReached(bool historicalBlockServingLimit) const EXCLUSIVE_LOCKS_REQUIRED(!m_total_bytes_sent_mutex)
check if the outbound target is reached if param historicalBlockServingLimit is set true,...
void StartExtraBlockRelayPeers()
void ForEachNode(const NodeFn &func) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
CSipHasher GetDeterministicRandomizer(uint64_t id) const
Get a unique deterministic randomizer.
bool EvictTxPeerIfFull(std::optional< NodeId > protect_peer=std::nullopt) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
If we are at capacity for inbound tx-relay peers, attempt to evict one.
uint32_t GetMappedAS(const CNetAddr &addr) const
bool MultipleManualOrFullOutboundConns(Network net) const EXCLUSIVE_LOCKS_REQUIRED(m_nodes_mutex)
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
enum Network GetNetwork() const
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 void UpdateLastBlockAnnounceTime(NodeId node, NodeClock::time_point time)=0
This function is used for testing the stale tip eviction logic, see denialofservice_tests....
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 CheckForStaleTipAndEvictPeers()=0
Evict extra outbound peers.
Store a list of transactions to be broadcast privately.
@ QueueFull
Rejected: the queue is already at MAX_TRANSACTIONS.
@ AlreadyPresent
The transaction was already present with send attempts remaining; no change.
@ Added
The transaction was newly added or reset after exhausting its send attempts.
static constexpr size_t MAX_TRANSACTIONS
Maximum number of transactions tracked simultaneously.
I randrange(I range) noexcept
Generate a random integer in the range [0..range), with range > 0.
bool Contains(Network net) const EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
std::string GetDebugMessage() const
std::string ToString() const
256-bit unsigned big integer.
constexpr bool IsNull() const
std::string ToString() const
CBlockIndex * LookupBlockIndex(const uint256 &hash) EXCLUSIVE_LOCKS_REQUIRED(cs_main)
CBlockFileInfo *GetBlockFileInfo(size_t n) EXCLUSIVE_LOCKS_REQUIRED(bool WriteBlockUndo(const CBlockUndo &blockundo, BlockValidationState &state, CBlockIndex &block) EXCLUSIVE_LOCKS_REQUIRED(FlatFilePos WriteBlock(const CBlock &block, int nHeight) EXCLUSIVE_LOCKS_REQUIRED(void UpdateBlockInfo(const CBlock &block, unsigned int nHeight, const FlatFilePos &pos) EXCLUSIVE_LOCKS_REQUIRED(bool IsPruneMode() const
Get block file info entry for one block file.
bool LoadingBlocks() const
ReadRawBlockResult ReadRawBlock(const FlatFilePos &pos, std::optional< std::pair< size_t, size_t > > block_part=std::nullopt) const
bool ReadBlock(CBlock &block, const FlatFilePos &pos, const std::optional< uint256 > &expected_hash) const
Functions for disk access for blocks.
Class responsible for deciding what transactions to request and, once downloaded, whether and how to ...
Manages warning messages within a node.
std::string ToString() const
const uint256 & ToUint256() const LIFETIMEBOUND
The util::Expected class provides a standard way for low-level functions to return either error value...
A token bucket rate limiter.
bool decrement(double n=1.0, double floor=0.0)
Consume n tokens.
void increment(const time_point &now)
Refill tokens based on elapsed time since last call.
double value() const
Current token balance.
The util::Unexpected class represents an unexpected value stored in util::Expected.
std::string TransportTypeAsString(TransportProtocolType transport_type)
Convert TransportProtocolType enum to a string value.
@ BLOCK_HEADER_LOW_WORK
the block header may be on a too-little-work chain
@ BLOCK_INVALID_HEADER
invalid proof of work or time too old
@ BLOCK_CACHED_INVALID
this block was cached as being invalid and we didn't store the reason why
@ BLOCK_CONSENSUS
invalid by consensus rules (excluding any below reasons)
@ BLOCK_MISSING_PREV
We don't have the previous block the checked one is built on.
@ BLOCK_INVALID_PREV
A block this one builds on is invalid.
@ BLOCK_MUTATED
the block's data didn't match the data committed to by the PoW
@ BLOCK_TIME_FUTURE
block timestamp was > 2 hours in the future (or our clock is bad)
@ BLOCK_RESULT_UNSET
initial value. Block has not yet been rejected
@ TX_MISSING_INPUTS
transaction was missing some of its inputs
@ TX_UNKNOWN
transaction was not validated because package failed
@ TX_NO_MEMPOOL
this node does not have a mempool so can't validate the transaction
@ TX_RESULT_UNSET
initial value. Tx has not yet been rejected
static size_t RecursiveDynamicUsage(const CScript &script)
RecursiveMutex cs_main
Mutex to guard access to validation specific variables, such as reading or changing the chainstate.
bool DeploymentActiveAfter(const CBlockIndex *pindexPrev, const Consensus::Params ¶ms, Consensus::BuriedDeployment dep, VersionBitsCache &versionbitscache)
Determine if a deployment is active for the next block.
bool DeploymentActiveAt(const CBlockIndex &index, const Consensus::Params ¶ms, Consensus::BuriedDeployment dep, VersionBitsCache &versionbitscache)
Determine if a deployment is active for this block.
is a home for simple enum and struct type definitions that can be used internally by functions in the...
#define LogDebug(category,...)
CSerializedNetMsg Make(std::string msg_type, Args &&... args)
constexpr const char * FILTERCLEAR
The filterclear message tells the receiving peer to remove a previously-set bloom filter.
constexpr const char * FEEFILTER
The feefilter message tells the receiving peer not to inv us any txs which do not meet the specified ...
constexpr const char * SENDHEADERS
Indicates that a node prefers to receive new block announcements via a "headers" message rather than ...
constexpr const char * GETBLOCKS
The getblocks message requests an inv message that provides block header hashes starting from a parti...
constexpr const char * HEADERS
The headers message sends one or more block headers to a node which previously requested certain head...
constexpr const char * ADDR
The addr (IP address) message relays connection information for peers on the network.
constexpr const char * GETBLOCKTXN
Contains a BlockTransactionsRequest Peer should respond with "blocktxn" message.
constexpr const char * CMPCTBLOCK
Contains a CBlockHeaderAndShortTxIDs object - providing a header and list of "short txids".
constexpr const char * CFCHECKPT
cfcheckpt is a response to a getcfcheckpt request containing a vector of evenly spaced filter headers...
constexpr const char * SENDADDRV2
The sendaddrv2 message signals support for receiving ADDRV2 messages (BIP155).
constexpr const char * GETADDR
The getaddr message requests an addr message from the receiving node, preferably one with lots of IP ...
constexpr const char * GETCFILTERS
getcfilters requests compact filters for a range of blocks.
constexpr const char * PONG
The pong message replies to a ping message, proving to the pinging node that the ponging node is stil...
constexpr const char * BLOCKTXN
Contains a BlockTransactions.
constexpr const char * CFHEADERS
cfheaders is a response to a getcfheaders request containing a filter header and a vector of filter h...
constexpr const char * PING
The ping message is sent periodically to help confirm that the receiving peer is still connected.
constexpr const char * FILTERLOAD
The filterload message tells the receiving peer to filter all relayed transactions and requested merk...
constexpr const char * SENDTXRCNCL
Contains a 4-byte version number and an 8-byte salt.
constexpr const char * ADDRV2
The addrv2 message relays connection information for peers on the network just like the addr message,...
constexpr const char * VERACK
The verack message acknowledges a previously-received version message, informing the connecting node ...
constexpr const char * GETHEADERS
The getheaders message requests a headers message that provides block headers starting from a particu...
constexpr const char * FILTERADD
The filteradd message tells the receiving peer to add a single element to a previously-set bloom filt...
constexpr const char * CFILTER
cfilter is a response to a getcfilters request containing a single compact filter.
constexpr const char * FEATURE
BIP 434 Peer feature negotiation.
constexpr const char * GETDATA
The getdata message requests one or more data objects from another node.
constexpr const char * SENDCMPCT
Contains a 1-byte bool and 8-byte LE version number.
constexpr const char * GETCFCHECKPT
getcfcheckpt requests evenly spaced compact filter headers, enabling parallelized download and valida...
constexpr const char * INV
The inv message (inventory message) transmits one or more inventories of objects known to the transmi...
constexpr const char * TX
The tx message transmits a single transaction.
constexpr const char * MEMPOOL
The mempool message requests the TXIDs of transactions that the receiving node has verified as valid ...
constexpr const char * NOTFOUND
The notfound message is a reply to a getdata message which requested an object the receiving node doe...
constexpr const char * MERKLEBLOCK
The merkleblock message is a reply to a getdata message which requested a block using the inventory t...
constexpr const char * WTXIDRELAY
Indicates that a node prefers to relay transactions via wtxid, rather than txid.
constexpr const char * BLOCK
The block message transmits a single serialized block.
constexpr const char * GETCFHEADERS
getcfheaders requests a compact filter header and the filter hashes for a range of blocks,...
constexpr const char * VERSION
The version message provides information about the transmitting node to the receiving node at the beg...
constexpr int32_t MAX_PEER_TX_ANNOUNCEMENTS
Maximum number of transactions to consider for requesting, per peer.
""_hex is a compile-time user-defined literal returning a std::array<std::byte>, equivalent to ParseH...
bool ShouldDebugLog(Category category)
Return whether messages with specified category should be debug logged.
std::string ToString(const T &t)
Locale-independent version of std::to_string.
std::string strSubVersion
Subversion as sent to the P2P network in version messages.
std::optional< CService > GetLocalAddrForPeer(CNode &node)
Returns a local address that we should advertise to this peer.
std::function< void(const CAddress &addr, const std::string &msg_type, std::span< const unsigned char > data, bool is_incoming)> CaptureMessage
Defaults to CaptureMessageToFile(), but can be overridden by unit tests.
bool SeenLocal(const CService &addr)
vote for a local address
constexpr unsigned int MAX_SUBVERSION_LENGTH
Maximum length of the user agent string in version message.
constexpr std::chrono::minutes TIMEOUT_INTERVAL
Time after which to disconnect, after waiting for a ping response (or inactivity).
static constexpr auto HEADERS_RESPONSE_TIME
How long to wait for a peer to respond to a getheaders request.
static constexpr size_t MAX_ADDR_TO_SEND
The maximum number of address records permitted in an ADDR message.
static constexpr auto INVENTORY_BUCKET_CHECK_DELAY
Delay between checking inventory bucket and backlog.
static constexpr size_t MAX_ADDR_PROCESSING_TOKEN_BUCKET
The soft limit of the address processing token bucket (the regular MAX_ADDR_RATE_PER_SECOND based inc...
static constexpr auto INVENTORY_BUCKET_BACKLOG_HEARTBEAT
Delay between inventory bucket backlog heartbeat log entries.
TRACEPOINT_SEMAPHORE(net, inbound_message)
static const int MAX_BLOCKS_IN_TRANSIT_PER_PEER
Number of blocks that can be requested at any given time from a single peer.
static constexpr auto BLOCK_STALLING_TIMEOUT_DEFAULT
Default time during which a peer must stall block download progress before being disconnected.
static constexpr auto AVG_FEEFILTER_BROADCAST_INTERVAL
Average delay between feefilter broadcasts in seconds.
static constexpr auto EXTRA_PEER_CHECK_INTERVAL
How frequently to check for extra outbound peers and disconnect.
static const unsigned int BLOCK_DOWNLOAD_WINDOW
Size of the "block download window": how far ahead of our current height do we fetch?...
static constexpr int STALE_RELAY_AGE_LIMIT
Age after which a stale block will no longer be served if requested as protection against fingerprint...
static constexpr int HISTORICAL_BLOCK_AGE
Age after which a block is considered historical for purposes of rate limiting block relay.
static constexpr auto ROTATE_ADDR_RELAY_DEST_INTERVAL
Delay between rotating the peers we relay a particular address to.
static constexpr auto MANUAL_PEER_BLOCK_DOWNLOAD_COOLDOWN
Time to avoid requesting blocks from a manual peer after it stalls block download.
static constexpr auto MINIMUM_CONNECT_TIME
Minimum time an outbound-peer-eviction candidate must be connected for, in order to evict.
static constexpr auto CHAIN_SYNC_TIMEOUT
Timeout for (unprotected) outbound peers to sync to our chainwork.
static constexpr auto OUTBOUND_INVENTORY_BROADCAST_INTERVAL
Average delay between trickled inventory transmissions for outbound peers.
static const unsigned int NODE_NETWORK_LIMITED_MIN_BLOCKS
Minimum blocks required to signal NODE_NETWORK_LIMITED.
static constexpr auto AVG_LOCAL_ADDRESS_BROADCAST_INTERVAL
Average delay between local address broadcasts.
static const int MAX_BLOCKTXN_DEPTH
Maximum depth of blocks we're willing to respond to GETBLOCKTXN requests for.
static constexpr size_t INVENTORY_BUCKET_BACKLOG_CAPACITY
Empty backlog target capacity.
static constexpr int32_t MAX_OUTBOUND_PEERS_TO_PROTECT_FROM_DISCONNECT
Protect at least this many outbound peers from disconnection due to slow/ behind headers chain.
static constexpr auto INBOUND_INVENTORY_BROADCAST_INTERVAL
Average delay between trickled inventory transmissions for inbound peers.
static constexpr size_t NUM_PRIVATE_BROADCAST_PER_TX
For private broadcast, send a transaction to this many peers.
static constexpr auto MAX_FEEFILTER_CHANGE_DELAY
Maximum feefilter broadcast delay after significant change.
static constexpr uint32_t MAX_GETCFILTERS_SIZE
Maximum number of compact filters that may be requested with one getcfilters.
static constexpr double OUTBOUND_INVENTORY_BUCKET_MULTIPLIER
Multiplier for the inventory bucket rate for outbounds.
static constexpr auto HEADERS_DOWNLOAD_TIMEOUT_BASE
Headers download timeout.
static const unsigned int MAX_GETDATA_SZ
Limit to avoid sending big packets.
static constexpr double BLOCK_DOWNLOAD_TIMEOUT_BASE
Block download timeout base, expressed in multiples of the block interval (i.e.
static constexpr auto PRIVATE_BROADCAST_MAX_CONNECTION_LIFETIME
Private broadcast connections must complete within this time.
static constexpr auto STALE_CHECK_INTERVAL
How frequently to check for stale tips.
static constexpr auto AVG_ADDRESS_BROADCAST_INTERVAL
Average delay between peer address broadcasts.
static const unsigned int MAX_LOCATOR_SZ
The maximum number of entries in a locator.
static constexpr double BLOCK_DOWNLOAD_TIMEOUT_PER_PEER
Additional block download timeout per parallel downloading peer (i.e.
static constexpr double MAX_ADDR_RATE_PER_SECOND
The maximum rate of address records we're willing to process on average.
static constexpr auto PING_INTERVAL
Time between pings automatically sent out for latency probing and keepalive.
static constexpr size_t INVENTORY_BUCKET_BACKLOG_HEARTBEAT_MIN
Minimum backlog to trigger heartbeat log entries.
static const int MAX_CMPCTBLOCK_DEPTH
Maximum depth of blocks we're willing to serve as compact blocks to peers when requested.
static const unsigned int MAX_BLOCKS_TO_ANNOUNCE
Maximum number of headers to announce when relaying blocks with headers message.
static const unsigned int NODE_NETWORK_LIMITED_ALLOW_CONN_BLOCKS
Window, in blocks, for connecting to NODE_NETWORK_LIMITED peers.
static constexpr uint32_t MAX_GETCFHEADERS_SIZE
Maximum number of cf hashes that may be requested with one getcfheaders.
static constexpr auto BLOCK_STALLING_TIMEOUT_MAX
Maximum timeout for stalling block download.
static constexpr auto HEADERS_DOWNLOAD_TIMEOUT_PER_HEADER
static constexpr uint64_t RANDOMIZER_ID_ADDRESS_RELAY
SHA256("main address relay")[0:8].
static constexpr size_t MAX_PCT_ADDR_TO_SEND
the maximum percentage of addresses from our addrman to return in response to a getaddr message.
static const unsigned int MAX_INV_SZ
The maximum number of entries in an 'inv' protocol message.
constexpr unsigned int MAX_CMPCTBLOCKS_INFLIGHT_PER_BLOCK
Maximum number of outstanding CMPCTBLOCK requests for the same block.
constexpr uint64_t CMPCTBLOCKS_VERSION
The compactblocks version we support.
ReachableNets g_reachable_nets
bool IsProxy(const CNetAddr &addr)
constexpr unsigned int DEFAULT_MIN_RELAY_TX_FEE
Default for -minrelaytxfee, minimum relay fee for transactions.
constexpr TransactionSerParams TX_NO_WITNESS
constexpr TransactionSerParams TX_WITH_WITNESS
std::shared_ptr< const CTransaction > CTransactionRef
GenTxid ToGenTxid(const CInv &inv)
Convert a TX/WITNESS_TX/WTX CInv to a GenTxid.
constexpr size_t MAX_FEATUREDATA_LENGTH
constexpr uint32_t MSG_WITNESS_FLAG
getdata message type flags
@ MSG_WTX
Defined in BIP 339.
@ MSG_CMPCT_BLOCK
Defined in BIP152.
@ MSG_WITNESS_BLOCK
Defined in BIP144.
ServiceFlags
nServices flags
constexpr size_t MAX_FEATUREID_LENGTH
static bool MayHaveUsefulAddressDB(ServiceFlags services)
Checks if a peer with the given service flags may be capable of having a robust address-storage DB.
constexpr int MIN_PEER_PROTO_VERSION
disconnect from peers older than this proto version
constexpr int SHORT_IDS_BLOCKS_VERSION
short-id-based block download starts with this version
constexpr int BIP0031_VERSION
BIP 0031, pong message, is enabled for all versions AFTER this one.
constexpr int FEEFILTER_VERSION
"feefilter" tells peers to filter invs to you by fee starts with this version
constexpr int WTXID_RELAY_VERSION
"wtxidrelay" message type for wtxid-based relay starts with this version
constexpr int INVALID_CB_NO_BAN_VERSION
not banning for invalid compact blocks starts with this version
constexpr int FEATURE_VERSION
"feature" message type for feature negotiation starts with this version
constexpr int SENDHEADERS_VERSION
"sendheaders" message type and announcing blocks with headers starts with this version
constexpr unsigned int MAX_SCRIPT_ELEMENT_SIZE
#define LIMITED_VECTOR(obj, n)
#define LIMITED_STRING(obj, n)
uint64_t ReadCompactSize(Stream &is, bool range_check=true)
Decode a CompactSize-encoded variable-length integer.
constexpr auto MakeUCharSpan(const V &v) -> decltype(UCharSpanCast(std::span{v}))
Like the std::span constructor, but for (const) unsigned char member types only.
Describes a place in the block chain to another node such that if the other node doesn't have the sam...
std::vector< uint256 > vHave
NodeClock::time_point m_last_block_announcement
NodeClock::duration m_ping_wait
std::vector< int > vHeightInFlight
CAmount m_fee_filter_received
std::chrono::seconds time_offset
bool m_addr_relay_enabled
uint64_t m_addr_rate_limited
uint64_t m_addr_processed
ServiceFlags their_services
Parameters that influence chain consensus.
int64_t nPowTargetSpacing
std::chrono::seconds PowTargetSpacing() const
Validation result for a transaction evaluated by MemPoolAccept (single or package).
const ResultType m_result_type
Result type.
const TxValidationState m_state
Contains information about why the transaction failed.
@ DIFFERENT_WITNESS
Valid, transaction was already in the mempool.
@ INVALID
Fully validated, valid.
const std::list< CTransactionRef > m_replaced_transactions
Mempool transactions replaced by the tx.
Version of the system clock that is mockable in the context of tests (via FakeNodeClock or SetMockTim...
static time_point now() noexcept
Return current system time or mocked time, if set.
std::chrono::time_point< NodeClock > time_point
static constexpr time_point epoch
Validation result for package mempool acceptance.
PackageValidationState m_state
std::map< Wtxid, MempoolAcceptResult > m_tx_results
Map from wtxid to finished MempoolAcceptResults.
std::chrono::seconds median_outbound_time_offset
Information about chainstate that notifications are sent from.
bool historical
Whether this is a historical chainstate downloading old blocks to validate an assumeutxo snapshot,...
CFeeRate min_relay_feerate
A fee rate smaller than this is considered zero fee (for relaying, mining and transaction creation)
std::vector< NodeId > m_senders
std::string ToString() const
#define AssertLockNotHeld(cs)
#define WITH_LOCK(cs, code)
Run code while locking a mutex.
COutPoint ProcessBlock(const NodeContext &node, const std::shared_ptr< CBlock > &block)
Returns the generated coin (or Null if the block was invalid).
#define EXCLUSIVE_LOCKS_REQUIRED(...)
#define LOCKS_EXCLUDED(...)
#define ACQUIRED_BEFORE(...)
#define TRACEPOINT(context,...)
consteval auto _(util::TranslatedLiteral str)
ReconciliationRegisterResult
constexpr uint32_t TXRECONCILIATION_VERSION
Supported transaction reconciliation protocol version.
std::string SanitizeString(std::string_view str, int rule)
Remove unsafe chars.
constexpr int64_t count_seconds(std::chrono::seconds t)
std::chrono::time_point< NodeClock, std::chrono::seconds > NodeSeconds
PackageMempoolAcceptResult ProcessNewPackage(Chainstate &active_chainstate, CTxMemPool &pool, const Package &package, bool test_accept, const std::optional< CFeeRate > &client_maxfeerate)
Validate (and maybe submit) a package to the mempool.
bool IsBlockMutated(const CBlock &block, bool check_witness_root)
Check if a block has been mutated (with respect to its merkle root and witness commitments).
bool HasValidProofOfWork(std::span< const CBlockHeader > headers, const Consensus::Params &consensusParams)
Check that the proof of work on each blockheader matches the value in nBits.
arith_uint256 CalculateClaimedHeadersWork(std::span< const CBlockHeader > headers)
Return the sum of the claimed work on a given set of headers.
@ UNVALIDATED
Blocks after an assumeutxo snapshot have been validated but the snapshot itself has not been validate...
constexpr unsigned int MIN_BLOCKS_TO_KEEP
Block files containing a block-height within MIN_BLOCKS_TO_KEEP of ActiveChain().Tip() will not be pr...