75#include <initializer_list>
212 std::unique_ptr<PartiallyDownloadedBlock> partialBlock;
246 std::atomic<ServiceFlags> m_their_services{
NODE_NONE};
249 const bool m_is_inbound;
252 Mutex m_misbehavior_mutex;
254 bool m_should_discourage
GUARDED_BY(m_misbehavior_mutex){
false};
257 Mutex m_block_inv_mutex;
261 std::vector<uint256> m_blocks_for_inv_relay
GUARDED_BY(m_block_inv_mutex);
265 std::vector<uint256> m_blocks_for_headers_relay
GUARDED_BY(m_block_inv_mutex);
276 std::atomic<uint64_t> m_ping_nonce_sent{0};
280 std::atomic<bool> m_ping_queued{
false};
283 std::atomic<bool> m_wtxid_relay{
false};
295 bool m_relay_txs
GUARDED_BY(m_bloom_filter_mutex){
false};
297 std::unique_ptr<CBloomFilter> m_bloom_filter
PT_GUARDED_BY(m_bloom_filter_mutex)
GUARDED_BY(m_bloom_filter_mutex){
nullptr};
308 std::vector<Wtxid> m_tx_inventory_to_send
GUARDED_BY(m_tx_inventory_mutex);
312 bool m_send_mempool
GUARDED_BY(m_tx_inventory_mutex){
false};
315 std::chrono::microseconds m_next_inv_send_time
GUARDED_BY(m_tx_inventory_mutex){0};
318 uint64_t m_last_inv_sequence
GUARDED_BY(m_tx_inventory_mutex){1};
321 std::atomic<CAmount> m_fee_filter_received{0};
327 LOCK(m_tx_relay_mutex);
329 m_tx_relay = std::make_unique<Peer::TxRelay>();
330 return m_tx_relay.get();
335 return WITH_LOCK(m_tx_relay_mutex,
return m_tx_relay.get());
364 std::atomic_bool m_addr_relay_enabled{
false};
368 mutable Mutex m_addr_send_times_mutex;
370 std::chrono::microseconds m_next_addr_send
GUARDED_BY(m_addr_send_times_mutex){0};
372 std::chrono::microseconds m_next_local_addr_send
GUARDED_BY(m_addr_send_times_mutex){0};
375 std::atomic_bool m_wants_addrv2{
false};
384 std::atomic<uint64_t> m_addr_rate_limited{0};
386 std::atomic<uint64_t> m_addr_processed{0};
392 Mutex m_getdata_requests_mutex;
394 std::deque<CInv> m_getdata_requests
GUARDED_BY(m_getdata_requests_mutex);
400 Mutex m_headers_sync_mutex;
403 std::unique_ptr<HeadersSyncState> m_headers_sync
PT_GUARDED_BY(m_headers_sync_mutex)
GUARDED_BY(m_headers_sync_mutex) {};
406 std::atomic<bool> m_sent_sendheaders{
false};
416 std::atomic<std::chrono::seconds> m_time_offset{0
s};
420 , m_our_services{our_services}
421 , m_is_inbound{is_inbound}
425 mutable Mutex m_tx_relay_mutex;
428 std::unique_ptr<TxRelay> m_tx_relay
GUARDED_BY(m_tx_relay_mutex);
431using PeerRef = std::shared_ptr<Peer>;
443 uint256 hashLastUnknownBlock{};
449 bool fSyncStarted{
false};
451 std::chrono::microseconds m_stalling_since{0us};
452 std::list<QueuedBlock> vBlocksInFlight;
454 std::chrono::microseconds m_downloading_since{0us};
456 bool fPreferredDownload{
false};
458 bool m_requested_hb_cmpctblocks{
false};
460 bool m_provides_cmpctblocks{
false};
486 struct ChainSyncTimeoutState {
488 std::chrono::seconds m_timeout{0
s};
492 bool m_sent_getheaders{
false};
494 bool m_protect{
false};
497 ChainSyncTimeoutState m_chain_sync;
500 int64_t m_last_block_announcement{0};
503struct InvToSendBucket {
504 const double count_floor{0};
505 std::vector<Wtxid> backlog;
522 static constexpr double SIZE_INIT{12'000'000};
523 static constexpr double SIZE_CAP{50'000'000};
524 static constexpr double SIZE_REFILL{20'000};
526 static constexpr double INBOUND_COUNT_SECONDS{30};
528 InvToSendBucket(
unsigned int rate,
double mult)
530 size_bucket(SIZE_REFILL * mult, SIZE_INIT, SIZE_CAP),
531 count_bucket(rate * mult, rate * INBOUND_COUNT_SECONDS, rate * INBOUND_COUNT_SECONDS)
537 return !backlog.empty() && size_bucket.
value() > 0 && count_bucket.
value() > 0;
548 bool decrement(
double size)
550 bool size_ok = size_bucket.
decrement(size, -50e3);
551 bool count_ok = count_bucket.
decrement(1, count_floor);
552 return size_ok && count_ok;
559 .count_bucket = count_bucket.
value(),
560 .size_bucket = size_bucket.
value(),
591 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);
593 EXCLUSIVE_LOCKS_REQUIRED(!m_peer_mutex, !m_most_recent_block_mutex, g_msgproc_mutex, !m_tx_download_mutex, !m_inv_to_send_mutex);
608 void SetBestBlock(
int height,
std::chrono::seconds time)
override
610 m_best_height = height;
611 m_best_block_time = time;
619 const std::atomic<bool>& interruptMsgProc)
620 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);
632 void ReattemptPrivateBroadcast(
CScheduler& scheduler);
647 void Misbehaving(Peer& peer, const
std::
string& message);
658 bool via_compact_block, const
std::
string& message = "")
667 bool MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer);
680 bool MaybeDisconnectForTxRelayCapacity(
CNode&
node, const
std::
string& msg_type,
695 bool first_time_failure)
720 bool ProcessOrphanTx(Peer& peer)
730 void ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
732 bool via_compact_block)
763 bool IsContinuationOfLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
777 bool TryLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
792 void HeadersDirectFetchBlocks(
CNode& pfrom, const Peer& peer, const
CBlockIndex& last_header);
794 void UpdatePeerStateForReceivedHeaders(
CNode& pfrom, Peer& peer, const
CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
801 template <
typename... Args>
802 void MakeAndPushMessage(
CNode&
node, std::string msg_type, Args&&...
args)
const
806 template <
typename... Args>
807 void MakeAndPushFeature(
CNode&
node, std::string_view feature_id, Args&&...
args)
const
810 std::vector<unsigned char> feature_data;
817 void PushNodeVersion(
CNode& pnode,
const Peer& peer);
866 std::unique_ptr<TxReconciliationTracker> m_txreconciliation;
869 std::atomic<int> m_best_height{-1};
871 std::atomic<std::chrono::seconds> m_best_block_time{0
s};
879 const Options m_opts;
881 bool RejectIncomingTxs(
const CNode& peer)
const;
889 mutable Mutex m_peer_mutex;
896 std::map<NodeId, PeerRef> m_peer_map
GUARDED_BY(m_peer_mutex);
906 uint32_t GetFetchFlags(
const Peer& peer)
const;
908 std::map<uint64_t, std::chrono::microseconds> m_next_inv_to_inbounds_per_network_key
GUARDED_BY(g_msgproc_mutex);
925 std::atomic<int> m_wtxid_relay_peers{0};
943 std::chrono::microseconds NextInvToInbounds(std::chrono::microseconds now,
944 std::chrono::seconds average_interval,
949 Mutex m_most_recent_block_mutex;
950 std::shared_ptr<const CBlock> m_most_recent_block
GUARDED_BY(m_most_recent_block_mutex);
951 std::shared_ptr<const CBlockHeaderAndShortTxIDs> m_most_recent_compact_block
GUARDED_BY(m_most_recent_block_mutex);
953 std::unique_ptr<const std::map<GenTxid, CTransactionRef>> m_most_recent_block_txs
GUARDED_BY(m_most_recent_block_mutex);
957 Mutex m_headers_presync_mutex;
965 using HeadersPresyncStats = std::pair<arith_uint256, std::optional<std::pair<int64_t, uint32_t>>>;
967 std::map<NodeId, HeadersPresyncStats> m_headers_presync_stats
GUARDED_BY(m_headers_presync_mutex) {};
971 std::atomic_bool m_headers_presync_should_signal{
false};
1041 std::atomic<
std::chrono::seconds> m_last_tip_update{0
s};
1047 void ProcessGetData(
CNode& pfrom, Peer& peer,
const std::atomic<bool>& interruptMsgProc)
1052 void ProcessBlock(
CNode&
node,
const std::shared_ptr<const CBlock>& block,
bool force_processing,
bool min_pow_checked);
1085 std::vector<std::pair<Wtxid, CTransactionRef>> vExtraTxnForCompact
GUARDED_BY(g_msgproc_mutex);
1087 size_t vExtraTxnForCompactIt
GUARDED_BY(g_msgproc_mutex) = 0;
1099 int64_t ApproximateBestBlockDepth() const;
1109 void ProcessGetBlockData(
CNode& pfrom, Peer& peer, const
CInv& inv)
1127 bool PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
1129 const
uint256& stop_hash, uint32_t max_height_diff,
1176 void ProcessAddrs(
std::string_view msg_type,
CNode& pfrom, Peer& peer,
std::vector<
CAddress>&& vAddr, const
std::atomic<
bool>& interruptMsgProc)
1182 void LogBlockHeader(const
CBlockIndex& index, const
CNode& peer,
bool via_compact_block);
1188 InvToSendBucket m_inbound_inv_bucket
GUARDED_BY(m_inv_to_send_mutex);
1189 InvToSendBucket m_outbound_inv_bucket
GUARDED_BY(m_inv_to_send_mutex);
1190 std::atomic<
NodeClock::time_point> m_next_inv_bucket_check{NodeClock::time_point::min()};
1191 std::optional<NodeClock::time_point> m_next_inv_bucket_heartbeat
GUARDED_BY(m_inv_to_send_mutex);
1196const CNodeState* PeerManagerImpl::
State(
NodeId pnode)
const
1198 std::map<NodeId, CNodeState>::const_iterator it = m_node_states.find(pnode);
1199 if (it == m_node_states.end())
1206 return const_cast<CNodeState*
>(std::as_const(*this).State(pnode));
1214static bool IsAddrCompatible(
const Peer& peer,
const CAddress& addr)
1219void PeerManagerImpl::AddAddressKnown(Peer& peer,
const CAddress& addr)
1221 assert(peer.m_addr_known);
1222 peer.m_addr_known->insert(addr.
GetKey());
1225void PeerManagerImpl::PushAddress(Peer& peer,
const CAddress& addr)
1230 assert(peer.m_addr_known);
1231 if (addr.
IsValid() && !peer.m_addr_known->contains(addr.
GetKey()) && IsAddrCompatible(peer, addr)) {
1233 peer.m_addrs_to_send[m_rng.randrange(peer.m_addrs_to_send.size())] = addr;
1235 peer.m_addrs_to_send.push_back(addr);
1240static void AddKnownTx(Peer& peer,
const uint256& hash)
1242 auto tx_relay = peer.GetTxRelay();
1243 if (!tx_relay)
return;
1245 LOCK(tx_relay->m_tx_inventory_mutex);
1246 tx_relay->m_tx_inventory_known_filter.insert(hash);
1250static bool CanServeBlocks(
const Peer& peer)
1257static bool IsLimitedPeer(
const Peer& peer)
1264static bool CanServeWitnesses(
const Peer& peer)
1269std::chrono::microseconds PeerManagerImpl::NextInvToInbounds(std::chrono::microseconds now,
1270 std::chrono::seconds average_interval,
1271 uint64_t network_key)
1273 auto [it, inserted] = m_next_inv_to_inbounds_per_network_key.try_emplace(network_key, 0us);
1274 auto& timer{it->second};
1276 timer = now + m_rng.rand_exp_duration(average_interval);
1281bool PeerManagerImpl::IsBlockRequested(
const uint256& hash)
1283 return mapBlocksInFlight.contains(hash);
1286bool PeerManagerImpl::IsBlockRequestedFromOutbound(
const uint256& hash)
1288 for (
auto range = mapBlocksInFlight.equal_range(hash); range.first != range.second; range.first++) {
1289 auto [nodeid, block_it] = range.first->second;
1290 PeerRef peer{GetPeerRef(nodeid)};
1291 if (peer && !peer->m_is_inbound)
return true;
1297void PeerManagerImpl::RemoveBlockRequest(
const uint256& hash, std::optional<NodeId> from_peer)
1299 auto range = mapBlocksInFlight.equal_range(hash);
1300 if (range.first == range.second) {
1308 while (range.first != range.second) {
1309 const auto& [node_id, list_it]{range.first->second};
1311 if (from_peer && *from_peer != node_id) {
1318 if (state.vBlocksInFlight.begin() == list_it) {
1320 state.m_downloading_since = std::max(state.m_downloading_since, GetTime<std::chrono::microseconds>());
1322 state.vBlocksInFlight.erase(list_it);
1324 if (state.vBlocksInFlight.empty()) {
1326 m_peers_downloading_from--;
1328 state.m_stalling_since = 0us;
1330 range.first = mapBlocksInFlight.erase(range.first);
1334bool PeerManagerImpl::BlockRequested(
NodeId nodeid,
const CBlockIndex& block, std::list<QueuedBlock>::iterator** pit)
1338 CNodeState *state =
State(nodeid);
1339 assert(state !=
nullptr);
1344 for (
auto range = mapBlocksInFlight.equal_range(hash); range.first != range.second; range.first++) {
1345 if (range.first->second.first == nodeid) {
1347 *pit = &range.first->second.second;
1354 RemoveBlockRequest(hash, nodeid);
1356 std::list<QueuedBlock>::iterator it = state->vBlocksInFlight.insert(state->vBlocksInFlight.end(),
1357 {&block, std::unique_ptr<PartiallyDownloadedBlock>(pit ? new PartiallyDownloadedBlock(&m_mempool) : nullptr)});
1358 if (state->vBlocksInFlight.size() == 1) {
1360 state->m_downloading_since = GetTime<std::chrono::microseconds>();
1361 m_peers_downloading_from++;
1363 auto itInFlight = mapBlocksInFlight.insert(std::make_pair(hash, std::make_pair(nodeid, it)));
1365 *pit = &itInFlight->second.second;
1370void PeerManagerImpl::MaybeSetPeerAsAnnouncingHeaderAndIDs(
NodeId nodeid)
1377 if (m_opts.ignore_incoming_txs)
return;
1379 CNodeState* nodestate =
State(nodeid);
1380 PeerRef peer{GetPeerRef(nodeid)};
1381 if (!nodestate || !nodestate->m_provides_cmpctblocks) {
1386 int num_outbound_hb_peers = 0;
1387 for (std::list<NodeId>::iterator it = lNodesAnnouncingHeaderAndIDs.begin(); it != lNodesAnnouncingHeaderAndIDs.end(); it++) {
1388 if (*it == nodeid) {
1389 lNodesAnnouncingHeaderAndIDs.erase(it);
1390 lNodesAnnouncingHeaderAndIDs.push_back(nodeid);
1393 PeerRef peer_ref{GetPeerRef(*it)};
1394 if (peer_ref && !peer_ref->m_is_inbound) ++num_outbound_hb_peers;
1396 if (peer && peer->m_is_inbound) {
1399 if (lNodesAnnouncingHeaderAndIDs.size() >= 3 && num_outbound_hb_peers == 1) {
1400 PeerRef remove_peer{GetPeerRef(lNodesAnnouncingHeaderAndIDs.front())};
1401 if (remove_peer && !remove_peer->m_is_inbound) {
1404 std::swap(lNodesAnnouncingHeaderAndIDs.front(), *std::next(lNodesAnnouncingHeaderAndIDs.begin()));
1413 lNodesAnnouncingHeaderAndIDs.push_back(pfrom->
GetId());
1416 if (nodeid_was_appended && lNodesAnnouncingHeaderAndIDs.size() > 3) {
1419 m_connman.
ForNode(lNodesAnnouncingHeaderAndIDs.front(), [
this](
CNode* pnodeStop) {
1422 pnodeStop->m_bip152_highbandwidth_to =
false;
1425 lNodesAnnouncingHeaderAndIDs.pop_front();
1429bool PeerManagerImpl::TipMayBeStale()
1433 if (m_last_tip_update.load() == 0
s) {
1434 m_last_tip_update = GetTime<std::chrono::seconds>();
1436 return m_last_tip_update.load() < GetTime<std::chrono::seconds>() - std::chrono::seconds{consensusParams.
nPowTargetSpacing * 3} && mapBlocksInFlight.empty();
1439int64_t PeerManagerImpl::ApproximateBestBlockDepth()
const
1444bool PeerManagerImpl::CanDirectFetch()
1451 if (state->pindexBestKnownBlock && pindex == state->pindexBestKnownBlock->GetAncestor(pindex->nHeight))
1453 if (state->pindexBestHeaderSent && pindex == state->pindexBestHeaderSent->GetAncestor(pindex->nHeight))
1458void PeerManagerImpl::ProcessBlockAvailability(
NodeId nodeid) {
1459 CNodeState *state =
State(nodeid);
1460 assert(state !=
nullptr);
1462 if (!state->hashLastUnknownBlock.IsNull()) {
1465 if (state->pindexBestKnownBlock ==
nullptr || pindex->
nChainWork >= state->pindexBestKnownBlock->nChainWork) {
1466 state->pindexBestKnownBlock = pindex;
1468 state->hashLastUnknownBlock.SetNull();
1473void PeerManagerImpl::UpdateBlockAvailability(
NodeId nodeid,
const uint256 &hash) {
1474 CNodeState *state =
State(nodeid);
1475 assert(state !=
nullptr);
1477 ProcessBlockAvailability(nodeid);
1482 if (state->pindexBestKnownBlock ==
nullptr || pindex->
nChainWork >= state->pindexBestKnownBlock->nChainWork) {
1483 state->pindexBestKnownBlock = pindex;
1487 state->hashLastUnknownBlock = hash;
1492void PeerManagerImpl::FindNextBlocksToDownload(
const Peer& peer,
unsigned int count, std::vector<const CBlockIndex*>& vBlocks,
NodeId& nodeStaller)
1497 vBlocks.reserve(vBlocks.size() +
count);
1498 CNodeState *state =
State(peer.m_id);
1499 assert(state !=
nullptr);
1502 ProcessBlockAvailability(peer.m_id);
1504 if (state->pindexBestKnownBlock ==
nullptr || state->pindexBestKnownBlock->nChainWork < m_chainman.
ActiveChain().
Tip()->
nChainWork || state->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
1515 state->pindexBestKnownBlock->GetAncestor(snap_base->nHeight) != snap_base) {
1516 LogDebug(
BCLog::NET,
"Not downloading blocks from peer=%d, which doesn't have the snapshot block in its best chain.\n", peer.m_id);
1525 if (state->pindexLastCommonBlock ==
nullptr ||
1526 fork_point->nChainWork > state->pindexLastCommonBlock->nChainWork ||
1527 state->pindexBestKnownBlock->GetAncestor(state->pindexLastCommonBlock->nHeight) != state->pindexLastCommonBlock) {
1528 state->pindexLastCommonBlock = fork_point;
1530 if (state->pindexLastCommonBlock == state->pindexBestKnownBlock)
1533 const CBlockIndex *pindexWalk = state->pindexLastCommonBlock;
1539 FindNextBlocks(vBlocks, peer, state, pindexWalk,
count, nWindowEnd, &m_chainman.
ActiveChain(), &nodeStaller);
1542void PeerManagerImpl::TryDownloadingHistoricalBlocks(
const Peer& peer,
unsigned int count, std::vector<const CBlockIndex*>& vBlocks,
const CBlockIndex *from_tip,
const CBlockIndex* target_block)
1547 if (vBlocks.size() >=
count) {
1551 vBlocks.reserve(
count);
1554 if (state->pindexBestKnownBlock ==
nullptr || state->pindexBestKnownBlock->GetAncestor(target_block->
nHeight) != target_block) {
1571void 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)
1573 std::vector<const CBlockIndex*> vToFetch;
1574 int nMaxHeight = std::min<int>(state->pindexBestKnownBlock->nHeight, nWindowEnd + 1);
1575 bool is_limited_peer = IsLimitedPeer(peer);
1577 while (pindexWalk->
nHeight < nMaxHeight) {
1581 int nToFetch = std::min(nMaxHeight - pindexWalk->
nHeight, std::max<int>(
count - vBlocks.size(), 128));
1582 vToFetch.resize(nToFetch);
1583 pindexWalk = state->pindexBestKnownBlock->
GetAncestor(pindexWalk->
nHeight + nToFetch);
1584 vToFetch[nToFetch - 1] = pindexWalk;
1585 for (
unsigned int i = nToFetch - 1; i > 0; i--) {
1586 vToFetch[i - 1] = vToFetch[i]->
pprev;
1606 state->pindexLastCommonBlock = pindex;
1613 if (waitingfor == -1) {
1615 waitingfor = mapBlocksInFlight.lower_bound(pindex->
GetBlockHash())->second.first;
1621 if (pindex->
nHeight > nWindowEnd) {
1623 if (vBlocks.size() == 0 && waitingfor != peer.m_id) {
1625 if (nodeStaller) *nodeStaller = waitingfor;
1635 vBlocks.push_back(pindex);
1636 if (vBlocks.size() ==
count) {
1645void PeerManagerImpl::PushNodeVersion(
CNode& pnode,
const Peer& peer)
1647 uint64_t my_services;
1649 uint64_t your_services;
1651 std::string my_user_agent;
1659 my_user_agent =
"/pynode:0.0.1/";
1661 my_tx_relay =
false;
1664 my_services = peer.m_our_services;
1665 my_time = TicksSinceEpoch<std::chrono::seconds>(
NodeClock::now());
1669 my_height = m_best_height;
1670 my_tx_relay = !RejectIncomingTxs(pnode);
1689 BCLog::NET,
"send version message: version=%d, blocks=%d%s, txrelay=%d, peer=%d\n",
1692 my_tx_relay, pnode.
GetId());
1695void PeerManagerImpl::UpdateLastBlockAnnounceTime(
NodeId node, int64_t time_in_seconds)
1699 if (state) state->m_last_block_announcement = time_in_seconds;
1707 m_node_states.try_emplace(m_node_states.end(), nodeid);
1709 WITH_LOCK(m_tx_download_mutex, m_txdownloadman.CheckIsEmpty(nodeid));
1715 PeerRef peer = std::make_shared<Peer>(nodeid, our_services,
node.IsInboundConn());
1718 m_peer_map.emplace_hint(m_peer_map.end(), nodeid, peer);
1722void PeerManagerImpl::ReattemptInitialBroadcast(
CScheduler& scheduler)
1726 for (
const auto& txid : unbroadcast_txids) {
1729 if (tx !=
nullptr) {
1730 InitiateTxBroadcastToAll(tx->GetWitnessHash());
1739 scheduler.
scheduleFromNow([&] { ReattemptInitialBroadcast(scheduler); }, delta);
1742void PeerManagerImpl::ReattemptPrivateBroadcast(
CScheduler& scheduler)
1746 size_t num_for_rebroadcast{0};
1747 const auto stale_txs = m_tx_for_private_broadcast.GetStale();
1748 if (!stale_txs.empty()) {
1749 for (
const auto& stale_tx : stale_txs) {
1755 "Reattempting broadcast of stale txid=%s wtxid=%s",
1756 stale_tx->GetHash().ToString(), stale_tx->GetWitnessHash().ToString());
1757 ++num_for_rebroadcast;
1760 stale_tx->GetHash().ToString(), stale_tx->GetWitnessHash().ToString(),
1761 mempool_acceptable.m_state.ToString());
1762 m_tx_for_private_broadcast.Remove(stale_tx);
1771 scheduler.
scheduleFromNow([&] { ReattemptPrivateBroadcast(scheduler); }, delta);
1774void PeerManagerImpl::FinalizeNode(
const CNode&
node)
1785 PeerRef peer = RemovePeer(nodeid);
1787 m_wtxid_relay_peers -= peer->m_wtxid_relay;
1788 assert(m_wtxid_relay_peers >= 0);
1790 CNodeState *state =
State(nodeid);
1791 assert(state !=
nullptr);
1793 if (state->fSyncStarted)
1796 for (
const QueuedBlock& entry : state->vBlocksInFlight) {
1797 auto range = mapBlocksInFlight.equal_range(entry.pindex->GetBlockHash());
1798 while (range.first != range.second) {
1799 auto [node_id, list_it] = range.first->second;
1800 if (node_id != nodeid) {
1803 range.first = mapBlocksInFlight.erase(range.first);
1808 LOCK(m_tx_download_mutex);
1809 m_txdownloadman.DisconnectedPeer(nodeid);
1811 if (m_txreconciliation) m_txreconciliation->ForgetPeer(nodeid);
1812 m_num_preferred_download_peers -= state->fPreferredDownload;
1813 m_peers_downloading_from -= (!state->vBlocksInFlight.empty());
1814 assert(m_peers_downloading_from >= 0);
1815 m_outbound_peers_with_protect_from_disconnect -= state->m_chain_sync.m_protect;
1816 assert(m_outbound_peers_with_protect_from_disconnect >= 0);
1818 m_node_states.erase(nodeid);
1820 if (m_node_states.empty()) {
1822 assert(mapBlocksInFlight.empty());
1823 assert(m_num_preferred_download_peers == 0);
1824 assert(m_peers_downloading_from == 0);
1825 assert(m_outbound_peers_with_protect_from_disconnect == 0);
1826 assert(m_wtxid_relay_peers == 0);
1827 WITH_LOCK(m_tx_download_mutex, m_txdownloadman.CheckIsEmpty());
1830 if (
node.fSuccessfullyConnected &&
1831 !
node.IsBlockOnlyConn() && !
node.IsPrivateBroadcastConn() && !
node.IsInboundConn()) {
1839 LOCK(m_headers_presync_mutex);
1840 m_headers_presync_stats.erase(nodeid);
1842 if (
node.IsPrivateBroadcastConn() &&
1843 !m_tx_for_private_broadcast.DidNodeConfirmReception(nodeid) &&
1844 m_tx_for_private_broadcast.HavePendingTransactions()) {
1851bool PeerManagerImpl::HasAllDesirableServiceFlags(
ServiceFlags services)
const
1854 return !(GetDesirableServiceFlags(services) & (~services));
1868PeerRef PeerManagerImpl::GetPeerRef(
NodeId id)
const
1871 auto it = m_peer_map.find(
id);
1872 return it != m_peer_map.end() ? it->second :
nullptr;
1875PeerRef PeerManagerImpl::RemovePeer(
NodeId id)
1879 auto it = m_peer_map.find(
id);
1880 if (it != m_peer_map.end()) {
1881 ret = std::move(it->second);
1882 m_peer_map.erase(it);
1887std::vector<PeerRef> PeerManagerImpl::GetAllPeers()
const
1889 std::vector<PeerRef> peers;
1891 peers.reserve(m_peer_map.size());
1892 for (
const auto& [
_, peer] : m_peer_map) {
1893 peers.push_back(peer);
1902 const CNodeState* state =
State(nodeid);
1903 if (state ==
nullptr)
1905 stats.
nSyncHeight = state->pindexBestKnownBlock ? state->pindexBestKnownBlock->nHeight : -1;
1906 stats.
nCommonHeight = state->pindexLastCommonBlock ? state->pindexLastCommonBlock->nHeight : -1;
1907 for (
const QueuedBlock& queue : state->vBlocksInFlight) {
1913 PeerRef peer = GetPeerRef(nodeid);
1914 if (peer ==
nullptr)
return false;
1922 NodeClock::duration ping_wait{0us};
1923 if ((0 != peer->m_ping_nonce_sent) && (peer->m_ping_start.load() >
NodeClock::epoch)) {
1927 if (
auto tx_relay = peer->GetTxRelay(); tx_relay !=
nullptr) {
1930 LOCK(tx_relay->m_tx_inventory_mutex);
1932 stats.
m_inv_to_send = tx_relay->m_tx_inventory_to_send.size();
1944 LOCK(peer->m_headers_sync_mutex);
1945 if (peer->m_headers_sync) {
1954std::vector<node::TxOrphanage::OrphanInfo> PeerManagerImpl::GetOrphanTransactions()
1956 LOCK(m_tx_download_mutex);
1957 return m_txdownloadman.GetOrphanTransactions();
1962 LOCK(m_inv_to_send_mutex);
1965 .ignores_incoming_txs = m_opts.ignore_incoming_txs,
1966 .private_broadcast = m_opts.private_broadcast,
1967 .tx_send_rate = m_opts.tx_send_rate,
1968 .inbound_bucket = m_inbound_inv_bucket.info(),
1969 .outbound_bucket = m_outbound_inv_bucket.info(),
1973std::vector<PrivateBroadcast::TxBroadcastInfo> PeerManagerImpl::GetPrivateBroadcastInfo()
const
1975 return m_tx_for_private_broadcast.GetBroadcastInfo();
1978std::vector<CTransactionRef> PeerManagerImpl::AbortPrivateBroadcast(
const uint256&
id)
1980 const auto snapshot{m_tx_for_private_broadcast.GetBroadcastInfo()};
1981 std::vector<CTransactionRef> removed_txs;
1983 size_t connections_cancelled{0};
1984 for (
const auto& tx_info : snapshot) {
1986 if (tx->GetHash().ToUint256() !=
id && tx->GetWitnessHash().ToUint256() !=
id)
continue;
1987 if (
const auto peer_acks{m_tx_for_private_broadcast.Remove(tx)}) {
1988 removed_txs.push_back(tx);
1999void PeerManagerImpl::AddToCompactExtraTransactions(
const CTransactionRef& tx)
2001 if (m_opts.max_extra_txs == 0)
return;
2002 if (vExtraTxnForCompact.size() < m_opts.max_extra_txs) {
2003 if (vExtraTxnForCompact.empty()) vExtraTxnForCompact.reserve(m_opts.max_extra_txs);
2004 vExtraTxnForCompact.emplace_back(tx->GetWitnessHash(), tx);
2006 vExtraTxnForCompact[vExtraTxnForCompactIt] = std::make_pair(tx->GetWitnessHash(), tx);
2008 vExtraTxnForCompactIt = (vExtraTxnForCompactIt + 1) % m_opts.max_extra_txs;
2011void PeerManagerImpl::Misbehaving(Peer& peer,
const std::string& message)
2013 LOCK(peer.m_misbehavior_mutex);
2015 const std::string message_prefixed = message.empty() ?
"" : (
": " + message);
2016 peer.m_should_discourage =
true;
2025 bool via_compact_block,
const std::string& message)
2027 PeerRef peer{GetPeerRef(nodeid)};
2038 if (!via_compact_block) {
2039 if (peer) Misbehaving(*peer, message);
2047 if (peer && !via_compact_block && !peer->m_is_inbound) {
2048 if (peer) Misbehaving(*peer, message);
2055 if (peer) Misbehaving(*peer, message);
2059 if (peer) Misbehaving(*peer, message);
2064 if (message !=
"") {
2069bool PeerManagerImpl::BlockRequestAllowed(
const CBlockIndex& block_index)
2089 PeerRef peer = GetPeerRef(peer_id);
2096 RemoveBlockRequest(block_index.
GetBlockHash(), std::nullopt);
2099 if (!BlockRequested(peer_id, block_index))
return util::Unexpected{
"Already requested from this peer"};
2122 return std::make_unique<PeerManagerImpl>(connman, addrman, banman, chainman, pool, warnings, opts);
2128 : m_rng{opts.deterministic_rng},
2130 m_chainparams(chainman.GetParams()),
2134 m_chainman(chainman),
2136 m_txdownloadman(
node::TxDownloadOptions{pool, m_rng, opts.deterministic_rng}),
2137 m_warnings{warnings},
2139 m_inbound_inv_bucket(m_opts.tx_send_rate, 1.0),
2144 if (opts.reconcile_txs) {
2149void PeerManagerImpl::StartScheduledTasks(
CScheduler& scheduler)
2160 scheduler.
scheduleFromNow([&] { ReattemptInitialBroadcast(scheduler); }, delta);
2162 if (m_opts.private_broadcast) {
2163 scheduler.
scheduleFromNow([&] { ReattemptPrivateBroadcast(scheduler); }, 0min);
2167void PeerManagerImpl::ActiveTipChange(
const CBlockIndex& new_tip,
bool is_ibd)
2175 LOCK(m_tx_download_mutex);
2179 m_txdownloadman.ActiveTipChange();
2189void PeerManagerImpl::BlockConnected(
2191 const std::shared_ptr<const CBlock>& pblock,
2196 m_last_tip_update = GetTime<std::chrono::seconds>();
2199 auto stalling_timeout = m_block_stalling_timeout.load();
2203 if (m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
2212 LOCK(m_tx_download_mutex);
2213 m_txdownloadman.BlockConnected(pblock);
2217void PeerManagerImpl::BlockDisconnected(
const std::shared_ptr<const CBlock> &block,
const CBlockIndex* pindex)
2219 LOCK(m_tx_download_mutex);
2220 m_txdownloadman.BlockDisconnected();
2227void PeerManagerImpl::NewPoWValidBlock(
const CBlockIndex *pindex,
const std::shared_ptr<const CBlock>& pblock)
2229 auto pcmpctblock = std::make_shared<const CBlockHeaderAndShortTxIDs>(*pblock,
FastRandomContext().rand64());
2233 if (pindex->
nHeight <= m_highest_fast_announce)
2235 m_highest_fast_announce = pindex->
nHeight;
2239 uint256 hashBlock(pblock->GetHash());
2240 const std::shared_future<CSerializedNetMsg> lazy_ser{
2244 auto most_recent_block_txs = std::make_unique<std::map<GenTxid, CTransactionRef>>();
2245 for (
const auto& tx : pblock->vtx) {
2246 most_recent_block_txs->emplace(tx->GetHash(), tx);
2247 most_recent_block_txs->emplace(tx->GetWitnessHash(), tx);
2250 LOCK(m_most_recent_block_mutex);
2251 m_most_recent_block_hash = hashBlock;
2252 m_most_recent_block = pblock;
2253 m_most_recent_compact_block = pcmpctblock;
2254 m_most_recent_block_txs = std::move(most_recent_block_txs);
2262 ProcessBlockAvailability(pnode->
GetId());
2266 if (state.m_requested_hb_cmpctblocks && !PeerHasHeader(&state, pindex) && PeerHasHeader(&state, pindex->
pprev)) {
2268 LogDebug(
BCLog::NET,
"%s sending header-and-ids %s to peer=%d\n",
"PeerManager::NewPoWValidBlock",
2269 hashBlock.ToString(), pnode->
GetId());
2272 PushMessage(*pnode, ser_cmpctblock.Copy());
2273 state.pindexBestHeaderSent = pindex;
2282void PeerManagerImpl::UpdatedBlockTip(
const CBlockIndex *pindexNew,
const CBlockIndex *pindexFork,
bool fInitialDownload)
2284 SetBestBlock(pindexNew->
nHeight, std::chrono::seconds{pindexNew->GetBlockTime()});
2287 if (fInitialDownload)
return;
2290 std::vector<uint256> vHashes;
2292 while (pindexToAnnounce != pindexFork) {
2294 pindexToAnnounce = pindexToAnnounce->
pprev;
2304 for (
auto& it : m_peer_map) {
2305 Peer& peer = *it.second;
2306 LOCK(peer.m_block_inv_mutex);
2307 for (
const uint256& hash : vHashes | std::views::reverse) {
2308 peer.m_blocks_for_headers_relay.push_back(hash);
2320void PeerManagerImpl::BlockChecked(
const std::shared_ptr<const CBlock>& block,
const BlockValidationState& state)
2324 const uint256 hash(block->GetHash());
2325 std::map<uint256, std::pair<NodeId, bool>>::iterator it = mapBlockSource.find(hash);
2330 it != mapBlockSource.end() &&
2331 State(it->second.first)) {
2332 MaybePunishNodeForBlock( it->second.first, state, !it->second.second);
2342 mapBlocksInFlight.count(hash) == mapBlocksInFlight.size()) {
2343 if (it != mapBlockSource.end()) {
2344 MaybeSetPeerAsAnnouncingHeaderAndIDs(it->second.first);
2347 if (it != mapBlockSource.end())
2348 mapBlockSource.erase(it);
2356bool PeerManagerImpl::AlreadyHaveBlock(
const uint256& block_hash)
2361void PeerManagerImpl::SendPings()
2364 for(
auto& it : m_peer_map) it.second->m_ping_queued =
true;
2367std::vector<Wtxid> InvToSendBucket::TakeForProcessing(
CTxMemPool& mempool)
2371 size_t n_to_take =
static_cast<size_t>(std::max<double>(count_bucket.
value() - count_floor, 0));
2373 std::vector<Wtxid> best;
2376 bool tokens_left =
true;
2377 for (
auto txiter : itervec) {
2378 auto& wtxid = txiter->GetTx().GetWitnessHash();
2380 best.push_back(wtxid);
2381 if (!decrement(txiter->GetTx().ComputeTotalSize())) {
2382 tokens_left =
false;
2385 backlog.push_back(wtxid);
2391 std::vector<Wtxid> dummy;
2393 dummy.swap(backlog);
2403 if (!backlog_bumped && now <= m_next_inv_bucket_check.load())
return;
2406 LOCK(m_inv_to_send_mutex);
2407 m_inbound_inv_bucket.increment(now);
2408 m_outbound_inv_bucket.increment(now);
2411 if (!m_next_inv_bucket_heartbeat.has_value()) {
2413 m_next_inv_bucket_heartbeat = now;
2416 if (m_next_inv_bucket_heartbeat.has_value() && now >= *m_next_inv_bucket_heartbeat) {
2417 LogDebug(
BCLog::NET,
"Transaction rate-limiting backlog inbound=%d itok=%.1f isz=%.1f outbound=%d otok=%.1f osz=%.1f",
2418 m_inbound_inv_bucket.backlog.size(),
2419 m_inbound_inv_bucket.count_bucket.value(),
2420 m_inbound_inv_bucket.size_bucket.value(),
2421 m_outbound_inv_bucket.backlog.size(),
2422 m_outbound_inv_bucket.count_bucket.value(),
2423 m_outbound_inv_bucket.size_bucket.value());
2424 if (m_inbound_inv_bucket.backlog.empty() && m_outbound_inv_bucket.backlog.empty()) {
2425 m_next_inv_bucket_heartbeat = std::nullopt;
2432 bool in_avail = m_inbound_inv_bucket.avail();
2433 bool out_avail = m_outbound_inv_bucket.avail();
2434 if (!in_avail && !out_avail)
return;
2436 std::vector<Wtxid> for_inbound;
2437 std::vector<Wtxid> for_outbound;
2441 if (in_avail) for_inbound = m_inbound_inv_bucket.TakeForProcessing(m_mempool);
2442 if (out_avail) for_outbound = m_outbound_inv_bucket.TakeForProcessing(m_mempool);
2445 if (!for_inbound.empty() || !for_outbound.empty()) {
2446 bool any_inbound_connected =
false;
2447 bool any_outbound_connected =
false;
2448 for (
const PeerRef& peer_ref : GetAllPeers()) {
2449 if (!peer_ref)
continue;
2450 Peer& peer{*peer_ref};
2451 auto tx_relay = peer.GetTxRelay();
2452 if (!tx_relay)
continue;
2454 LOCK(tx_relay->m_tx_inventory_mutex);
2460 if (tx_relay->m_next_inv_send_time == 0
s)
continue;
2461 if (peer.m_is_inbound) {
2462 any_inbound_connected =
true;
2464 any_outbound_connected =
true;
2466 for (
auto& i : (peer.m_is_inbound ? for_inbound : for_outbound)) {
2467 tx_relay->m_tx_inventory_to_send.push_back(i);
2474 if (!any_inbound_connected) m_inbound_inv_bucket.backlog.clear();
2475 if (!any_outbound_connected) m_outbound_inv_bucket.backlog.clear();
2479void PeerManagerImpl::InitiateTxBroadcastToAll(
const Wtxid& wtxid)
2482 LOCK(m_inv_to_send_mutex);
2483 m_inbound_inv_bucket.backlog.push_back(wtxid);
2484 m_outbound_inv_bucket.backlog.push_back(wtxid);
2491 const auto txstr{
strprintf(
"txid=%s, wtxid=%s", tx->GetHash().ToString(), tx->GetWitnessHash().ToString())};
2492 switch (m_tx_for_private_broadcast.Add(tx)) {
2507void PeerManagerImpl::RelayAddress(
NodeId originator,
2523 const auto current_time{GetTime<std::chrono::seconds>()};
2531 unsigned int nRelayNodes = (fReachable || (hasher.Finalize() & 1)) ? 2 : 1;
2533 std::array<std::pair<uint64_t, Peer*>, 2> best{{{0,
nullptr}, {0,
nullptr}}};
2534 assert(nRelayNodes <= best.size());
2538 for (
auto& [
id, peer] : m_peer_map) {
2539 if (peer->m_addr_relay_enabled &&
id != originator && IsAddrCompatible(*peer, addr)) {
2541 for (
unsigned int i = 0; i < nRelayNodes; i++) {
2542 if (hashKey > best[i].first) {
2543 std::copy(best.begin() + i, best.begin() + nRelayNodes - 1, best.begin() + i + 1);
2544 best[i] = std::make_pair(hashKey, peer.get());
2551 for (
unsigned int i = 0; i < nRelayNodes && best[i].first != 0; i++) {
2552 PushAddress(*best[i].second, addr);
2556void PeerManagerImpl::ProcessGetBlockData(
CNode& pfrom, Peer& peer,
const CInv& inv)
2558 std::shared_ptr<const CBlock> a_recent_block;
2559 std::shared_ptr<const CBlockHeaderAndShortTxIDs> a_recent_compact_block;
2561 LOCK(m_most_recent_block_mutex);
2562 a_recent_block = m_most_recent_block;
2563 a_recent_compact_block = m_most_recent_compact_block;
2566 bool need_activate_chain =
false;
2578 need_activate_chain =
true;
2582 if (need_activate_chain) {
2584 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
2591 bool can_direct_fetch{
false};
2599 if (!BlockRequestAllowed(*pindex)) {
2600 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block that isn't in the main chain\n", __func__, pfrom.
GetId());
2627 can_direct_fetch = CanDirectFetch();
2631 std::shared_ptr<const CBlock> pblock;
2632 if (a_recent_block && a_recent_block->GetHash() == inv.
hash) {
2633 pblock = a_recent_block;
2651 std::shared_ptr<CBlock> pblockRead = std::make_shared<CBlock>();
2661 pblock = pblockRead;
2669 bool sendMerkleBlock =
false;
2671 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
2672 LOCK(tx_relay->m_bloom_filter_mutex);
2673 if (tx_relay->m_bloom_filter) {
2674 sendMerkleBlock =
true;
2675 merkleBlock =
CMerkleBlock(*pblock, *tx_relay->m_bloom_filter);
2678 if (sendMerkleBlock) {
2686 for (
const auto& [tx_idx,
_] : merkleBlock.
vMatchedTxn)
2697 if (a_recent_compact_block && a_recent_compact_block->header.GetHash() == inv.
hash) {
2710 LOCK(peer.m_block_inv_mutex);
2712 if (inv.
hash == peer.m_continuation_block) {
2716 std::vector<CInv> vInv;
2717 vInv.emplace_back(
MSG_BLOCK, tip->GetBlockHash());
2719 peer.m_continuation_block.SetNull();
2727 auto txinfo{std::visit(
2728 [&](
const auto&
id) {
2729 return m_mempool.
info_for_relay(
id,
WITH_LOCK(tx_relay.m_tx_inventory_mutex,
return tx_relay.m_last_inv_sequence));
2733 return std::move(txinfo.tx);
2738 LOCK(m_most_recent_block_mutex);
2739 if (m_most_recent_block_txs !=
nullptr) {
2740 auto it = m_most_recent_block_txs->find(gtxid);
2741 if (it != m_most_recent_block_txs->end())
return it->second;
2748void PeerManagerImpl::ProcessGetData(
CNode& pfrom, Peer& peer,
const std::atomic<bool>& interruptMsgProc)
2752 auto tx_relay = peer.GetTxRelay();
2754 std::deque<CInv>::iterator it = peer.m_getdata_requests.begin();
2755 std::vector<CInv> vNotFound;
2760 while (it != peer.m_getdata_requests.end() && it->IsGenTxMsg()) {
2761 if (interruptMsgProc)
return;
2766 const CInv &inv = *it++;
2768 if (tx_relay ==
nullptr) {
2774 if (
auto tx{FindTxForGetData(*tx_relay,
ToGenTxid(inv))}) {
2777 MakeAndPushMessage(pfrom,
NetMsgType::TX, maybe_with_witness(*tx));
2780 vNotFound.push_back(inv);
2786 if (it != peer.m_getdata_requests.end() && !pfrom.
fPauseSend) {
2787 const CInv &inv = *it++;
2789 ProcessGetBlockData(pfrom, peer, inv);
2798 peer.m_getdata_requests.erase(peer.m_getdata_requests.begin(), it);
2800 if (!vNotFound.empty()) {
2819uint32_t PeerManagerImpl::GetFetchFlags(
const Peer& peer)
const
2821 uint32_t nFetchFlags = 0;
2822 if (CanServeWitnesses(peer)) {
2831 for (
size_t i = 0; i < req.
indexes.size(); i++) {
2833 Misbehaving(peer,
"getblocktxn with out-of-bounds tx indices");
2840 uint32_t tx_requested_size{0};
2841 for (
const auto& tx : resp.txn) tx_requested_size += tx->ComputeTotalSize();
2847bool PeerManagerImpl::CheckHeadersPoW(
const std::vector<CBlockHeader>&
headers, Peer& peer)
2851 Misbehaving(peer,
"header with invalid proof of work");
2856 if (!CheckHeadersAreContinuous(
headers)) {
2857 Misbehaving(peer,
"non-continuous headers sequence");
2882void PeerManagerImpl::HandleUnconnectingHeaders(
CNode& pfrom, Peer& peer,
2883 const std::vector<CBlockHeader>&
headers)
2887 if (MaybeSendGetHeaders(pfrom,
GetLocator(best_header), peer)) {
2888 LogDebug(
BCLog::NET,
"received header %s: missing prev block %s, sending getheaders (%d) to end (peer=%d)\n",
2890 headers[0].hashPrevBlock.ToString(),
2891 best_header->nHeight,
2901bool PeerManagerImpl::CheckHeadersAreContinuous(
const std::vector<CBlockHeader>&
headers)
const
2905 if (!hashLastBlock.
IsNull() && header.hashPrevBlock != hashLastBlock) {
2908 hashLastBlock = header.GetHash();
2913bool PeerManagerImpl::IsContinuationOfLowWorkHeadersSync(Peer& peer,
CNode& pfrom, std::vector<CBlockHeader>&
headers)
2915 if (peer.m_headers_sync) {
2916 auto result = peer.m_headers_sync->ProcessNextHeaders(
headers,
headers.size() == m_opts.max_headers_result);
2918 if (result.success) peer.m_last_getheaders_timestamp = {};
2919 if (result.request_more) {
2920 auto locator = peer.m_headers_sync->NextHeadersRequestLocator();
2922 Assume(!locator.vHave.empty());
2925 if (!locator.vHave.empty()) {
2928 bool sent_getheaders = MaybeSendGetHeaders(pfrom, locator, peer);
2931 locator.vHave.front().ToString(), pfrom.
GetId());
2936 peer.m_headers_sync.reset(
nullptr);
2941 LOCK(m_headers_presync_mutex);
2942 m_headers_presync_stats.erase(pfrom.
GetId());
2945 HeadersPresyncStats stats;
2946 stats.first = peer.m_headers_sync->GetPresyncWork();
2948 stats.second = {peer.m_headers_sync->GetPresyncHeight(),
2949 peer.m_headers_sync->GetPresyncTime()};
2953 LOCK(m_headers_presync_mutex);
2954 m_headers_presync_stats[pfrom.
GetId()] = stats;
2955 auto best_it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
2956 bool best_updated =
false;
2957 if (best_it == m_headers_presync_stats.end()) {
2961 const HeadersPresyncStats* stat_best{
nullptr};
2962 for (
const auto& [peer, stat] : m_headers_presync_stats) {
2963 if (!stat_best || stat > *stat_best) {
2968 m_headers_presync_bestpeer = peer_best;
2969 best_updated = (peer_best == pfrom.
GetId());
2970 }
else if (best_it->first == pfrom.
GetId() || stats > best_it->second) {
2972 m_headers_presync_bestpeer = pfrom.
GetId();
2973 best_updated =
true;
2975 if (best_updated && stats.second.has_value()) {
2977 m_headers_presync_should_signal =
true;
2981 if (result.success) {
2984 headers.swap(result.pow_validated_headers);
2987 return result.success;
2995bool PeerManagerImpl::TryLowWorkHeadersSync(Peer& peer,
CNode& pfrom,
const CBlockIndex& chain_start_header, std::vector<CBlockHeader>&
headers)
3002 arith_uint256 minimum_chain_work = GetAntiDoSWorkThreshold();
3006 if (total_work < minimum_chain_work) {
3010 if (
headers.size() == m_opts.max_headers_result) {
3020 LOCK(peer.m_headers_sync_mutex);
3022 m_chainparams.
HeadersSync(), chain_start_header, minimum_chain_work));
3027 (void)IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3041bool PeerManagerImpl::IsAncestorOfBestHeaderOrTip(
const CBlockIndex* header)
3043 if (header ==
nullptr) {
3045 }
else if (m_chainman.m_best_header !=
nullptr && header == m_chainman.m_best_header->GetAncestor(header->
nHeight)) {
3053bool PeerManagerImpl::MaybeSendGetHeaders(
CNode& pfrom,
const CBlockLocator& locator, Peer& peer)
3061 peer.m_last_getheaders_timestamp = current_time;
3072void PeerManagerImpl::HeadersDirectFetchBlocks(
CNode& pfrom,
const Peer& peer,
const CBlockIndex& last_header)
3075 CNodeState *nodestate =
State(pfrom.
GetId());
3078 std::vector<const CBlockIndex*> vToFetch;
3086 vToFetch.push_back(pindexWalk);
3088 pindexWalk = pindexWalk->
pprev;
3100 std::vector<CInv> vGetData;
3102 for (
const CBlockIndex* pindex : vToFetch | std::views::reverse) {
3107 uint32_t nFetchFlags = GetFetchFlags(peer);
3109 BlockRequested(pfrom.
GetId(), *pindex);
3113 if (vGetData.size() > 1) {
3118 if (vGetData.size() > 0) {
3119 if (!m_opts.ignore_incoming_txs &&
3120 nodestate->m_provides_cmpctblocks &&
3121 vGetData.size() == 1 &&
3122 mapBlocksInFlight.size() == 1 &&
3138void PeerManagerImpl::UpdatePeerStateForReceivedHeaders(
CNode& pfrom, Peer& peer,
3139 const CBlockIndex& last_header,
bool received_new_header,
bool may_have_more_headers)
3142 CNodeState *nodestate =
State(pfrom.
GetId());
3151 nodestate->m_last_block_announcement =
GetTime();
3159 if (nodestate->pindexBestKnownBlock && nodestate->pindexBestKnownBlock->nChainWork < m_chainman.
MinimumChainWork()) {
3181 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) {
3183 nodestate->m_chain_sync.m_protect =
true;
3184 ++m_outbound_peers_with_protect_from_disconnect;
3189void PeerManagerImpl::ProcessHeadersMessage(
CNode& pfrom, Peer& peer,
3190 std::vector<CBlockHeader>&&
headers,
3191 bool via_compact_block)
3193 size_t nCount =
headers.size();
3200 LOCK(peer.m_headers_sync_mutex);
3201 if (peer.m_headers_sync) {
3202 peer.m_headers_sync.reset(
nullptr);
3203 LOCK(m_headers_presync_mutex);
3204 m_headers_presync_stats.erase(pfrom.
GetId());
3208 peer.m_last_getheaders_timestamp = {};
3216 if (!CheckHeadersPoW(
headers, peer)) {
3231 bool already_validated_work =
false;
3234 bool have_headers_sync =
false;
3236 LOCK(peer.m_headers_sync_mutex);
3238 already_validated_work = IsContinuationOfLowWorkHeadersSync(peer, pfrom,
headers);
3254 have_headers_sync = !!peer.m_headers_sync;
3259 bool headers_connect_blockindex{chain_start_header !=
nullptr};
3261 if (!headers_connect_blockindex) {
3265 HandleUnconnectingHeaders(pfrom, peer,
headers);
3272 peer.m_last_getheaders_timestamp = {};
3282 already_validated_work = already_validated_work || IsAncestorOfBestHeaderOrTip(last_received_header);
3289 already_validated_work =
true;
3295 if (!already_validated_work && TryLowWorkHeadersSync(peer, pfrom,
3296 *chain_start_header,
headers)) {
3308 bool received_new_header{last_received_header ==
nullptr};
3314 state, &pindexLast)};
3320 "If this happens with all peers, consider database corruption (that -reindex may fix) "
3321 "or a potential consensus incompatibility.",
3324 MaybePunishNodeForBlock(pfrom.
GetId(), state, via_compact_block,
"invalid header received");
3330 if (processed && received_new_header) {
3331 LogBlockHeader(*pindexLast, pfrom,
false);
3335 if (nCount == m_opts.max_headers_result && !have_headers_sync) {
3337 if (MaybeSendGetHeaders(pfrom,
GetLocator(pindexLast), peer)) {
3342 UpdatePeerStateForReceivedHeaders(pfrom, peer, *pindexLast, received_new_header, nCount == m_opts.max_headers_result);
3345 HeadersDirectFetchBlocks(pfrom, peer, *pindexLast);
3351 bool first_time_failure)
3357 PeerRef peer{GetPeerRef(nodeid)};
3360 ptx->GetHash().ToString(),
3361 ptx->GetWitnessHash().ToString(),
3365 const auto& [add_extra_compact_tx, unique_parents, package_to_validate] = m_txdownloadman.MempoolRejectedTx(ptx, state, nodeid, first_time_failure);
3368 AddToCompactExtraTransactions(ptx);
3370 for (
const Txid& parent_txid : unique_parents) {
3371 if (peer) AddKnownTx(*peer, parent_txid.ToUint256());
3374 return package_to_validate;
3377void PeerManagerImpl::ProcessValidTx(
NodeId nodeid,
const CTransactionRef& tx,
const std::list<CTransactionRef>& replaced_transactions)
3383 m_txdownloadman.MempoolAcceptedTx(tx);
3387 tx->GetHash().ToString(),
3388 tx->GetWitnessHash().ToString(),
3391 InitiateTxBroadcastToAll(tx->GetWitnessHash());
3394 AddToCompactExtraTransactions(removedTx);
3404 const auto&
package = package_to_validate.m_txns;
3405 const auto& senders = package_to_validate.
m_senders;
3408 m_txdownloadman.MempoolRejectedPackage(package);
3411 if (!
Assume(package.size() == 2))
return;
3415 auto package_iter = package.rbegin();
3416 auto senders_iter = senders.rbegin();
3417 while (package_iter != package.rend()) {
3418 const auto& tx = *package_iter;
3419 const NodeId nodeid = *senders_iter;
3420 const auto it_result{package_result.
m_tx_results.find(tx->GetWitnessHash())};
3424 const auto& tx_result = it_result->second;
3425 switch (tx_result.m_result_type) {
3428 ProcessValidTx(nodeid, tx, tx_result.m_replaced_transactions);
3438 ProcessInvalidTx(nodeid, tx, tx_result.m_state,
false);
3456bool PeerManagerImpl::ProcessOrphanTx(Peer& peer)
3463 while (
CTransactionRef porphanTx = m_txdownloadman.GetTxToReconsider(peer.m_id)) {
3466 const Txid& orphanHash = porphanTx->GetHash();
3467 const Wtxid& orphan_wtxid = porphanTx->GetWitnessHash();
3484 ProcessInvalidTx(peer.m_id, porphanTx, state,
false);
3493bool PeerManagerImpl::PrepareBlockFilterRequest(
CNode&
node, Peer& peer,
3495 const uint256& stop_hash, uint32_t max_height_diff,
3499 const bool supported_filter_type =
3502 if (!supported_filter_type) {
3504 static_cast<uint8_t
>(filter_type),
node.DisconnectMsg());
3505 node.fDisconnect =
true;
3514 if (!stop_index || !BlockRequestAllowed(*stop_index)) {
3517 node.fDisconnect =
true;
3522 uint32_t stop_height = stop_index->
nHeight;
3523 if (start_height > stop_height) {
3525 "start height %d and stop height %d, %s",
3526 start_height, stop_height,
node.DisconnectMsg());
3527 node.fDisconnect =
true;
3530 if (stop_height - start_height >= max_height_diff) {
3532 stop_height - start_height + 1, max_height_diff,
node.DisconnectMsg());
3533 node.fDisconnect =
true;
3538 if (!filter_index) {
3548 uint8_t filter_type_ser;
3549 uint32_t start_height;
3552 vRecv >> filter_type_ser >> start_height >> stop_hash;
3558 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3563 std::vector<BlockFilter> filters;
3565 LogDebug(
BCLog::NET,
"Failed to find block filter in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3570 for (
const auto& filter : filters) {
3577 uint8_t filter_type_ser;
3578 uint32_t start_height;
3581 vRecv >> filter_type_ser >> start_height >> stop_hash;
3587 if (!PrepareBlockFilterRequest(
node, peer, filter_type, start_height, stop_hash,
3593 if (start_height > 0) {
3595 stop_index->
GetAncestor(
static_cast<int>(start_height - 1));
3597 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3603 std::vector<uint256> filter_hashes;
3605 LogDebug(
BCLog::NET,
"Failed to find block filter hashes in index: filter_type=%s, start_height=%d, stop_hash=%s\n",
3619 uint8_t filter_type_ser;
3622 vRecv >> filter_type_ser >> stop_hash;
3628 if (!PrepareBlockFilterRequest(
node, peer, filter_type, 0, stop_hash,
3629 std::numeric_limits<uint32_t>::max(),
3630 stop_index, filter_index)) {
3638 for (
int i =
headers.size() - 1; i >= 0; i--) {
3643 LogDebug(
BCLog::NET,
"Failed to find block filter header in index: filter_type=%s, block_hash=%s\n",
3657 bool new_block{
false};
3658 m_chainman.
ProcessNewBlock(block, force_processing, min_pow_checked, &new_block);
3660 node.m_last_block_time = GetTime<std::chrono::seconds>();
3665 RemoveBlockRequest(block->GetHash(), std::nullopt);
3668 mapBlockSource.erase(block->GetHash());
3672void PeerManagerImpl::ProcessCompactBlockTxns(
CNode& pfrom, Peer& peer,
const BlockTransactions& block_transactions)
3674 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
3675 bool fBlockRead{
false};
3679 auto range_flight = mapBlocksInFlight.equal_range(block_transactions.
blockhash);
3680 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
3681 bool requested_block_from_this_peer{
false};
3684 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
3686 while (range_flight.first != range_flight.second) {
3687 auto [node_id, block_it] = range_flight.first->second;
3688 if (node_id == pfrom.
GetId() && block_it->partialBlock) {
3689 requested_block_from_this_peer =
true;
3692 range_flight.first++;
3695 if (!requested_block_from_this_peer) {
3707 Misbehaving(peer,
"previous compact block reconstruction attempt failed");
3718 Misbehaving(peer,
"invalid compact block/non-matching block transactions");
3721 if (first_in_flight) {
3726 std::vector<CInv> invs;
3731 LogDebug(
BCLog::NET,
"Peer %d sent us a compact block but it failed to reconstruct, waiting on first download to complete\n", pfrom.
GetId());
3744 mapBlockSource.emplace(block_transactions.
blockhash, std::make_pair(pfrom.
GetId(),
false));
3759void PeerManagerImpl::LogBlockHeader(
const CBlockIndex& index,
const CNode& peer,
bool via_compact_block) {
3771 "Saw new %sheader hash=%s height=%d %s",
3772 via_compact_block ?
"cmpctblock " :
"",
3784void PeerManagerImpl::PushPrivateBroadcastTx(
CNode&
node)
3788 const auto opt_tx{m_tx_for_private_broadcast.PickTxForSend(
node.GetId(),
CService{node.addr})};
3791 node.fDisconnect =
true;
3797 tx->GetHash().ToString(), tx->HasWitness() ?
strprintf(
", wtxid=%s", tx->GetWitnessHash().ToString()) :
"",
3803void PeerManagerImpl::ProcessMessage(Peer& peer,
CNode& pfrom,
const std::string& msg_type,
DataStream& vRecv,
3805 const std::atomic<bool>& interruptMsgProc)
3820 uint64_t nNonce = 1;
3823 std::string cleanSubVer;
3824 int starting_height = -1;
3827 vRecv >> nVersion >> Using<CustomUintFormatter<8>>(nServices) >> nTime;
3842 LogDebug(
BCLog::NET,
"peer does not offer the expected services (%08x offered, %08x expected), %s",
3844 GetDesirableServiceFlags(nServices),
3857 if (!vRecv.
empty()) {
3865 if (!vRecv.
empty()) {
3866 std::string strSubVer;
3870 if (!vRecv.
empty()) {
3871 vRecv >> starting_height;
3891 PushNodeVersion(pfrom, peer);
3895 const int greatest_common_version = std::min(nVersion, pfrom.
AdvertisedVersion());
3900 peer.m_their_services = nServices;
3904 pfrom.cleanSubVer = cleanSubVer;
3915 (fRelay || (peer.m_our_services &
NODE_BLOOM))) {
3916 auto*
const tx_relay = peer.SetTxRelay();
3918 LOCK(tx_relay->m_bloom_filter_mutex);
3919 tx_relay->m_relay_txs = fRelay;
3925 LogDebug(
BCLog::NET,
"receive version message: %s: version %d, blocks=%d, us=%s, txrelay=%d, %s%s",
3926 cleanSubVer.empty() ?
"<no user agent>" : cleanSubVer, pfrom.
nVersion,
3928 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
3946 if (greatest_common_version >= 70016) {
3961 const auto* tx_relay = peer.GetTxRelay();
3962 if (tx_relay &&
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs) &&
3964 const uint64_t recon_salt = m_txreconciliation->PreRegisterPeer(pfrom.
GetId());
3977 if (MaybeDisconnectForTxRelayCapacity(pfrom, msg_type, pfrom.
GetId()))
return;
3985 m_num_preferred_download_peers += state->fPreferredDownload;
3991 bool send_getaddr{
false};
3993 send_getaddr = SetupAddressRelay(pfrom, peer);
4003 peer.m_getaddr_sent =
true;
4027 peer.m_time_offset =
NodeSeconds{std::chrono::seconds{nTime}} - Now<NodeSeconds>();
4031 m_outbound_time_offsets.Add(peer.m_time_offset);
4032 m_outbound_time_offsets.WarnIfOutOfSync();
4036 if (greatest_common_version <= 70012) {
4037 constexpr auto finalAlert{
"60010000000000000000000000ffffff7f00000000ffffff7ffeffff7f01ffffff7f00000000ffffff7f00ffffff7f002f555247454e543a20416c657274206b657920636f6d70726f6d697365642c2075706772616465207265717569726564004630440220653febd6410f470f6bae11cad19c48413becb1ac2c17f908fd0fd53bdc3abd5202206d0e9c96fe88d4a0f01ed9dedae2b6f9e00da94cad0fecaae66ecf689bf71b50"_hex};
4038 MakeAndPushMessage(pfrom,
"alert", finalAlert);
4061 auto new_peer_msg = [&]() {
4063 return strprintf(
"New %s peer connected: transport: %s, version: %d, %s%s",
4067 (mapped_as ?
strprintf(
", mapped_as=%d", mapped_as) :
""));
4075 LogInfo(
"%s", new_peer_msg());
4078 if (
auto tx_relay = peer.GetTxRelay()) {
4087 tx_relay->m_tx_inventory_mutex,
4088 return tx_relay->m_tx_inventory_to_send.empty() &&
4089 tx_relay->m_next_inv_send_time == 0
s));
4099 PushPrivateBroadcastTx(pfrom);
4112 if (m_txreconciliation) {
4113 if (!peer.m_wtxid_relay || !m_txreconciliation->IsPeerRegistered(pfrom.
GetId())) {
4117 m_txreconciliation->ForgetPeer(pfrom.
GetId());
4123 const CNodeState* state =
State(pfrom.
GetId());
4125 .m_preferred = state->fPreferredDownload,
4126 .m_relay_permissions = pfrom.HasPermission(NetPermissionFlags::Relay),
4127 .m_wtxid_relay = peer.m_wtxid_relay,
4136 peer.m_prefers_headers =
true;
4141 uint8_t sendcmpct_hb{0};
4142 uint64_t sendcmpct_version{0};
4143 vRecv >> sendcmpct_hb >> sendcmpct_version;
4147 if (sendcmpct_hb > 1) {
4148 Misbehaving(peer,
"invalid sendcmpct announce field");
4156 CNodeState* nodestate =
State(pfrom.
GetId());
4157 nodestate->m_provides_cmpctblocks =
true;
4158 nodestate->m_requested_hb_cmpctblocks = sendcmpct_hb;
4175 if (!peer.m_wtxid_relay) {
4176 peer.m_wtxid_relay =
true;
4177 m_wtxid_relay_peers++;
4196 peer.m_wants_addrv2 =
true;
4213 std::string feature_id;
4217 std::vector<unsigned char> feature_data_vec;
4220 }
catch (
const std::exception&) {
4223 if (feature_id.size() < 4 || !vRecv.
empty()) {
4243 if (!m_txreconciliation) {
4244 LogDebug(
BCLog::NET,
"sendtxrcncl from peer=%d ignored, as our node does not have txreconciliation enabled\n", pfrom.
GetId());
4255 if (RejectIncomingTxs(pfrom)) {
4264 const auto* tx_relay = peer.GetTxRelay();
4265 if (!tx_relay || !
WITH_LOCK(tx_relay->m_bloom_filter_mutex,
return tx_relay->m_relay_txs)) {
4271 uint32_t peer_txreconcl_version;
4272 uint64_t remote_salt;
4273 vRecv >> peer_txreconcl_version >> remote_salt;
4276 peer_txreconcl_version, remote_salt);
4308 const auto ser_params{
4316 std::vector<CAddress> vAddr;
4317 vRecv >> ser_params(vAddr);
4318 ProcessAddrs(msg_type, pfrom, peer, std::move(vAddr), interruptMsgProc);
4323 std::vector<CInv> vInv;
4327 Misbehaving(peer,
strprintf(
"inv message size = %u", vInv.size()));
4331 const bool reject_tx_invs{RejectIncomingTxs(pfrom)};
4335 const auto current_time{GetTime<std::chrono::microseconds>()};
4338 for (
CInv& inv : vInv) {
4339 if (interruptMsgProc)
return;
4344 if (peer.m_wtxid_relay) {
4351 const bool fAlreadyHave = AlreadyHaveBlock(inv.
hash);
4354 UpdateBlockAvailability(pfrom.
GetId(), inv.
hash);
4362 best_block = &inv.
hash;
4365 if (reject_tx_invs) {
4371 AddKnownTx(peer, inv.
hash);
4374 const bool fAlreadyHave{m_txdownloadman.AddTxAnnouncement(pfrom.
GetId(), gtxid, current_time)};
4382 if (best_block !=
nullptr) {
4394 if (state.fSyncStarted || (!peer.m_inv_triggered_getheaders_before_sync && *best_block != m_last_block_inv_triggering_headers_sync)) {
4395 if (MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer)) {
4397 m_chainman.m_best_header->nHeight, best_block->ToString(),
4400 if (!state.fSyncStarted) {
4401 peer.m_inv_triggered_getheaders_before_sync =
true;
4405 m_last_block_inv_triggering_headers_sync = *best_block;
4414 std::vector<CInv> vInv;
4418 Misbehaving(peer,
strprintf(
"getdata message size = %u", vInv.size()));
4424 if (vInv.size() > 0) {
4429 const auto pushed_tx_opt{m_tx_for_private_broadcast.GetTxForNode(pfrom.
GetId())};
4430 if (!pushed_tx_opt) {
4441 if (vInv.size() == 1 && vInv[0].IsMsgTx() && vInv[0].hash == pushed_tx->GetHash().ToUint256()) {
4445 peer.m_ping_queued =
true;
4456 LOCK(peer.m_getdata_requests_mutex);
4457 peer.m_getdata_requests.insert(peer.m_getdata_requests.end(), vInv.begin(), vInv.end());
4458 ProcessGetData(pfrom, peer, interruptMsgProc);
4467 vRecv >> locator >> hashStop;
4483 std::shared_ptr<const CBlock> a_recent_block;
4485 LOCK(m_most_recent_block_mutex);
4486 a_recent_block = m_most_recent_block;
4489 if (!m_chainman.
ActiveChainstate().ActivateBestChain(state, a_recent_block)) {
4504 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4519 if (--nLimit <= 0) {
4523 WITH_LOCK(peer.m_block_inv_mutex, {peer.m_continuation_block = pindex->GetBlockHash();});
4535 for (
size_t i = 1; i < req.
indexes.size(); ++i) {
4539 std::shared_ptr<const CBlock> recent_block;
4541 LOCK(m_most_recent_block_mutex);
4542 if (m_most_recent_block_hash == req.
blockhash)
4543 recent_block = m_most_recent_block;
4547 SendBlockTransactions(pfrom, peer, *recent_block, req);
4566 if (!block_pos.IsNull()) {
4573 SendBlockTransactions(pfrom, peer, block, req);
4586 WITH_LOCK(peer.m_getdata_requests_mutex, peer.m_getdata_requests.push_back(inv));
4594 vRecv >> locator >> hashStop;
4612 if (m_chainman.
ActiveTip() ==
nullptr ||
4614 LogDebug(
BCLog::NET,
"Ignoring getheaders from peer=%d because active chain has too little work; sending empty response\n", pfrom.
GetId());
4621 CNodeState *nodestate =
State(pfrom.
GetId());
4630 if (!BlockRequestAllowed(*pindex)) {
4631 LogDebug(
BCLog::NET,
"%s: ignoring request from peer=%i for old block header that isn't in the main chain\n", __func__, pfrom.
GetId());
4644 std::vector<CBlock> vHeaders;
4645 int nLimit = m_opts.max_headers_result;
4647 for (; pindex; pindex = m_chainman.
ActiveChain().Next(*pindex))
4650 if (--nLimit <= 0 || pindex->GetBlockHash() == hashStop)
4665 nodestate->pindexBestHeaderSent = pindex ? pindex : m_chainman.
ActiveChain().
Tip();
4671 if (RejectIncomingTxs(pfrom)) {
4685 const Txid& txid = ptx->GetHash();
4686 const Wtxid& wtxid = ptx->GetWitnessHash();
4689 AddKnownTx(peer, hash);
4691 if (
const auto num_broadcasted{m_tx_for_private_broadcast.Remove(ptx)}) {
4693 "network from %s; stopping private broadcast attempts",
4704 const auto& [should_validate, package_to_validate] = m_txdownloadman.ReceivedTx(pfrom.
GetId(), ptx);
4705 if (!should_validate) {
4710 if (!m_mempool.
exists(txid)) {
4711 LogInfo(
"Not relaying non-mempool transaction %s (wtxid=%s) from forcerelay peer=%d\n",
4714 LogInfo(
"Force relaying tx %s (wtxid=%s) from peer=%d\n",
4716 InitiateTxBroadcastToAll(wtxid);
4720 if (package_to_validate) {
4723 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4724 ProcessPackageResult(package_to_validate.value(), package_result);
4730 Assume(!package_to_validate.has_value());
4740 if (
auto package_to_validate{ProcessInvalidTx(pfrom.
GetId(), ptx, state,
true)}) {
4743 package_result.
m_state.
IsValid() ?
"package accepted" :
"package rejected");
4744 ProcessPackageResult(package_to_validate.value(), package_result);
4757 }
else if (m_opts.ignore_incoming_txs) {
4764 const CNodeState *nodestate =
State(pfrom.
GetId());
4765 if (!nodestate->m_provides_cmpctblocks) {
4772 vRecv >> cmpctblock;
4774 bool received_new_header =
false;
4784 MaybeSendGetHeaders(pfrom,
GetLocator(m_chainman.m_best_header), peer);
4794 received_new_header =
true;
4802 MaybePunishNodeForBlock(pfrom.
GetId(), state,
true,
"invalid header via cmpctblock");
4809 if (received_new_header) {
4810 LogBlockHeader(*pindex, pfrom,
true);
4813 bool fProcessBLOCKTXN =
false;
4817 bool fRevertToHeaderProcessing =
false;
4821 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
4822 bool fBlockReconstructed =
false;
4828 CNodeState *nodestate =
State(pfrom.
GetId());
4833 nodestate->m_last_block_announcement =
GetTime();
4839 auto range_flight = mapBlocksInFlight.equal_range(pindex->
GetBlockHash());
4840 size_t already_in_flight = std::distance(range_flight.first, range_flight.second);
4841 bool requested_block_from_this_peer{
false};
4844 bool first_in_flight = already_in_flight == 0 || (range_flight.first->second.first == pfrom.
GetId());
4846 while (range_flight.first != range_flight.second) {
4847 if (range_flight.first->second.first == pfrom.
GetId()) {
4848 requested_block_from_this_peer =
true;
4851 range_flight.first++;
4861 if (requested_block_from_this_peer) {
4864 std::vector<CInv> vInv(1);
4872 if (!already_in_flight && !CanDirectFetch()) {
4880 requested_block_from_this_peer) {
4881 std::list<QueuedBlock>::iterator* queuedBlockIt =
nullptr;
4882 if (!BlockRequested(pfrom.
GetId(), *pindex, &queuedBlockIt)) {
4883 if (!(*queuedBlockIt)->partialBlock)
4896 Misbehaving(peer,
"invalid compact block");
4899 if (first_in_flight) {
4901 std::vector<CInv> vInv(1);
4912 for (
size_t i = 0; i < cmpctblock.
BlockTxCount(); i++) {
4917 fProcessBLOCKTXN =
true;
4918 }
else if (first_in_flight) {
4925 IsBlockRequestedFromOutbound(blockhash) ||
4944 ReadStatus status = tempBlock.InitData(cmpctblock, vExtraTxnForCompact);
4949 std::vector<CTransactionRef> dummy;
4951 status = tempBlock.FillBlock(*pblock, dummy,
4954 fBlockReconstructed =
true;
4958 if (requested_block_from_this_peer) {
4961 std::vector<CInv> vInv(1);
4967 fRevertToHeaderProcessing =
true;
4972 if (fProcessBLOCKTXN) {
4975 return ProcessCompactBlockTxns(pfrom, peer, txn);
4978 if (fRevertToHeaderProcessing) {
4984 return ProcessHeadersMessage(pfrom, peer, {cmpctblock.
header},
true);
4987 if (fBlockReconstructed) {
4992 mapBlockSource.emplace(pblock->GetHash(), std::make_pair(pfrom.
GetId(),
false));
5010 RemoveBlockRequest(pblock->GetHash(), std::nullopt);
5027 return ProcessCompactBlockTxns(pfrom, peer, resp);
5038 std::vector<CBlockHeader>
headers;
5042 if (nCount > m_opts.max_headers_result) {
5043 Misbehaving(peer,
strprintf(
"headers message size = %u", nCount));
5047 for (
unsigned int n = 0; n < nCount; n++) {
5052 ProcessHeadersMessage(pfrom, peer, std::move(
headers),
false);
5056 if (m_headers_presync_should_signal.exchange(
false)) {
5057 HeadersPresyncStats stats;
5059 LOCK(m_headers_presync_mutex);
5060 auto it = m_headers_presync_stats.find(m_headers_presync_bestpeer);
5061 if (it != m_headers_presync_stats.end()) stats = it->second;
5079 std::shared_ptr<CBlock> pblock = std::make_shared<CBlock>();
5090 Misbehaving(peer,
"mutated block");
5095 bool forceProcessing =
false;
5096 const uint256 hash(pblock->GetHash());
5097 bool min_pow_checked =
false;
5102 forceProcessing = IsBlockRequested(hash);
5103 RemoveBlockRequest(hash, pfrom.
GetId());
5107 mapBlockSource.emplace(hash, std::make_pair(pfrom.
GetId(),
true));
5111 min_pow_checked =
true;
5114 ProcessBlock(pfrom, pblock, forceProcessing, min_pow_checked);
5131 Assume(SetupAddressRelay(pfrom, peer));
5135 if (peer.m_getaddr_recvd) {
5139 peer.m_getaddr_recvd =
true;
5141 peer.m_addrs_to_send.clear();
5142 std::vector<CAddress> vAddr;
5148 for (
const CAddress &addr : vAddr) {
5149 PushAddress(peer, addr);
5177 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5178 LOCK(tx_relay->m_tx_inventory_mutex);
5179 tx_relay->m_send_mempool =
true;
5205 ProcessPong(pfrom, peer, time_received, vRecv);
5221 Misbehaving(peer,
"too-large bloom filter");
5222 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5224 LOCK(tx_relay->m_bloom_filter_mutex);
5225 tx_relay->m_bloom_filter.reset(
new CBloomFilter(filter));
5226 tx_relay->m_relay_txs =
true;
5230 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5241 std::vector<unsigned char> vData;
5249 }
else if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5250 LOCK(tx_relay->m_bloom_filter_mutex);
5251 if (tx_relay->m_bloom_filter) {
5252 tx_relay->m_bloom_filter->insert(vData);
5258 Misbehaving(peer,
"bad filteradd message");
5269 auto tx_relay = peer.GetTxRelay();
5270 if (!tx_relay)
return;
5273 LOCK(tx_relay->m_bloom_filter_mutex);
5274 tx_relay->m_bloom_filter =
nullptr;
5275 tx_relay->m_relay_txs =
true;
5279 MaybeDisconnectForTxRelayCapacity(pfrom, msg_type);
5285 vRecv >> newFeeFilter;
5287 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
5288 tx_relay->m_fee_filter_received = newFeeFilter;
5296 ProcessGetCFilters(pfrom, peer, vRecv);
5301 ProcessGetCFHeaders(pfrom, peer, vRecv);
5306 ProcessGetCFCheckPt(pfrom, peer, vRecv);
5311 std::vector<CInv> vInv;
5313 std::vector<GenTxid> tx_invs;
5315 for (
CInv &inv : vInv) {
5321 LOCK(m_tx_download_mutex);
5322 m_txdownloadman.ReceivedNotFound(pfrom.
GetId(), tx_invs);
5331bool PeerManagerImpl::MaybeDiscourageAndDisconnect(
CNode& pnode, Peer& peer)
5334 LOCK(peer.m_misbehavior_mutex);
5337 if (!peer.m_should_discourage)
return false;
5339 peer.m_should_discourage =
false;
5344 LogWarning(
"Not punishing noban peer %d!", peer.m_id);
5350 LogWarning(
"Not punishing manually connected peer %d!", peer.m_id);
5370bool PeerManagerImpl::MaybeDisconnectForTxRelayCapacity(
CNode&
node,
const std::string& msg_type, std::optional<NodeId> protect_peer)
5372 if (!
node.IsInboundConn() || !
node.m_relays_txs)
return false;
5375 LogDebug(
BCLog::NET,
"failed to find a tx-relaying eviction candidate - connection dropped after %s message, peer=%d\n", msg_type,
node.GetId());
5376 node.fDisconnect =
true;
5380bool PeerManagerImpl::ProcessMessages(
CNode&
node, std::atomic<bool>& interruptMsgProc)
5385 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
5386 if (maybe_peer ==
nullptr)
return false;
5387 Peer& peer{*maybe_peer};
5391 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent)
return false;
5394 LOCK(peer.m_getdata_requests_mutex);
5395 if (!peer.m_getdata_requests.empty()) {
5396 ProcessGetData(
node, peer, interruptMsgProc);
5400 const bool processed_orphan = ProcessOrphanTx(peer);
5402 if (
node.fDisconnect)
5405 if (processed_orphan)
return true;
5410 LOCK(peer.m_getdata_requests_mutex);
5411 if (!peer.m_getdata_requests.empty())
return true;
5415 if (
node.fPauseSend)
return false;
5417 auto poll_result{
node.PollMessage()};
5424 bool fMoreWork = poll_result->second;
5428 node.m_addr_name.c_str(),
5429 node.ConnectionTypeAsString().c_str(),
5435 if (m_opts.capture_messages) {
5440 ProcessMessage(peer,
node,
msg.m_type,
msg.m_recv,
msg.m_time, interruptMsgProc);
5441 if (interruptMsgProc)
return false;
5443 LOCK(peer.m_getdata_requests_mutex);
5444 if (!peer.m_getdata_requests.empty()) fMoreWork =
true;
5451 LOCK(m_tx_download_mutex);
5452 if (m_txdownloadman.HaveMoreWork(peer.m_id)) fMoreWork =
true;
5453 }
catch (
const std::exception& e) {
5462void PeerManagerImpl::ConsiderEviction(
CNode& pto, Peer& peer, std::chrono::seconds time_in_seconds)
5475 if (state.pindexBestKnownBlock !=
nullptr && state.pindexBestKnownBlock->nChainWork >= m_chainman.
ActiveChain().
Tip()->
nChainWork) {
5477 if (state.m_chain_sync.m_timeout != 0
s) {
5478 state.m_chain_sync.m_timeout = 0
s;
5479 state.m_chain_sync.m_work_header =
nullptr;
5480 state.m_chain_sync.m_sent_getheaders =
false;
5482 }
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)) {
5490 state.m_chain_sync.m_work_header = m_chainman.
ActiveChain().
Tip();
5491 state.m_chain_sync.m_sent_getheaders =
false;
5492 }
else if (state.m_chain_sync.m_timeout > 0
s && time_in_seconds > state.m_chain_sync.m_timeout) {
5496 if (state.m_chain_sync.m_sent_getheaders) {
5498 LogInfo(
"Outbound peer has old chain, best known block = %s, %s", state.pindexBestKnownBlock !=
nullptr ? state.pindexBestKnownBlock->GetBlockHash().ToString() :
"<none>", pto.
DisconnectMsg());
5501 assert(state.m_chain_sync.m_work_header);
5506 MaybeSendGetHeaders(pto,
5507 GetLocator(state.m_chain_sync.m_work_header->pprev),
5509 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());
5510 state.m_chain_sync.m_sent_getheaders =
true;
5531 std::pair<NodeId, std::chrono::seconds> youngest_peer{-1, 0}, next_youngest_peer{-1, 0};
5535 if (pnode->
GetId() > youngest_peer.first) {
5536 next_youngest_peer = youngest_peer;
5537 youngest_peer.first = pnode->GetId();
5538 youngest_peer.second = pnode->m_last_block_time;
5541 NodeId to_disconnect = youngest_peer.first;
5542 if (youngest_peer.second > next_youngest_peer.second) {
5545 to_disconnect = next_youngest_peer.first;
5554 CNodeState *node_state =
State(pnode->
GetId());
5555 if (node_state ==
nullptr ||
5558 LogDebug(
BCLog::NET,
"disconnecting extra block-relay-only peer=%d (last block received at time %d)\n",
5562 LogDebug(
BCLog::NET,
"keeping block-relay-only peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5563 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), node_state->vBlocksInFlight.size());
5578 int64_t oldest_block_announcement = std::numeric_limits<int64_t>::max();
5581 AssertLockHeld(::cs_main);
5585 if (!pnode->IsFullOutboundConn() || pnode->fDisconnect) return;
5586 CNodeState *state = State(pnode->GetId());
5587 if (state == nullptr) return;
5589 if (state->m_chain_sync.m_protect) return;
5592 if (!m_connman.MultipleManualOrFullOutboundConns(pnode->addr.GetNetwork())) return;
5593 if (state->m_last_block_announcement < oldest_block_announcement || (state->m_last_block_announcement == oldest_block_announcement && pnode->GetId() > worst_peer)) {
5594 worst_peer = pnode->GetId();
5595 oldest_block_announcement = state->m_last_block_announcement;
5598 if (worst_peer != -1) {
5609 LogDebug(
BCLog::NET,
"disconnecting extra outbound peer=%d (last block announcement received at time %d)\n", pnode->
GetId(), oldest_block_announcement);
5613 LogDebug(
BCLog::NET,
"keeping outbound peer=%d chosen for eviction (connect time: %d, blocks_in_flight: %d)\n",
5614 pnode->
GetId(), TicksSinceEpoch<std::chrono::seconds>(pnode->
m_connected), state.vBlocksInFlight.size());
5630void PeerManagerImpl::CheckForStaleTipAndEvictPeers()
5635 auto now{GetTime<std::chrono::seconds>()};
5637 EvictExtraOutboundPeers(current_time);
5639 if (now > m_stale_tip_check_time) {
5643 LogInfo(
"Potential stale tip detected, will try using extra outbound peer (last tip update: %d seconds ago)\n",
5652 if (!m_initial_sync_finished && CanDirectFetch()) {
5654 m_initial_sync_finished =
true;
5661 peer.m_ping_nonce_sent &&
5671 bool pingSend =
false;
5673 if (peer.m_ping_queued) {
5678 if (peer.m_ping_nonce_sent == 0 && now > peer.m_ping_start.load() +
PING_INTERVAL) {
5687 }
while (
nonce == 0);
5688 peer.m_ping_queued =
false;
5689 peer.m_ping_start = now;
5691 peer.m_ping_nonce_sent =
nonce;
5695 peer.m_ping_nonce_sent = 0;
5701void PeerManagerImpl::MaybeSendAddr(
CNode&
node, Peer& peer, std::chrono::microseconds current_time)
5704 if (!peer.m_addr_relay_enabled)
return;
5706 LOCK(peer.m_addr_send_times_mutex);
5709 peer.m_next_local_addr_send < current_time) {
5716 if (peer.m_next_local_addr_send != 0us) {
5717 peer.m_addr_known->reset();
5720 CAddress local_addr{*local_service, peer.m_our_services, Now<NodeSeconds>()};
5721 if (peer.m_next_local_addr_send == 0us) {
5725 if (IsAddrCompatible(peer, local_addr)) {
5726 std::vector<CAddress> self_announcement{local_addr};
5727 if (peer.m_wants_addrv2) {
5735 PushAddress(peer, local_addr);
5742 if (current_time <= peer.m_next_addr_send)
return;
5755 bool ret = peer.m_addr_known->contains(addr.
GetKey());
5756 if (!
ret) peer.m_addr_known->insert(addr.
GetKey());
5759 peer.m_addrs_to_send.erase(std::remove_if(peer.m_addrs_to_send.begin(), peer.m_addrs_to_send.end(), addr_already_known),
5760 peer.m_addrs_to_send.end());
5763 if (peer.m_addrs_to_send.empty())
return;
5765 if (peer.m_wants_addrv2) {
5770 peer.m_addrs_to_send.clear();
5773 if (peer.m_addrs_to_send.capacity() > 40) {
5774 peer.m_addrs_to_send.shrink_to_fit();
5778void PeerManagerImpl::MaybeSendSendHeaders(
CNode&
node, Peer& peer)
5786 CNodeState &state = *
State(
node.GetId());
5787 if (state.pindexBestKnownBlock !=
nullptr &&
5794 peer.m_sent_sendheaders =
true;
5799void PeerManagerImpl::MaybeSendFeefilter(
CNode& pto, Peer& peer, std::chrono::microseconds current_time)
5801 if (m_opts.ignore_incoming_txs)
return;
5817 if (peer.m_fee_filter_sent == MAX_FILTER) {
5820 peer.m_next_send_feefilter = 0us;
5823 if (current_time > peer.m_next_send_feefilter) {
5824 CAmount filterToSend = m_fee_filter_rounder.round(currentFilter);
5827 if (filterToSend != peer.m_fee_filter_sent) {
5829 peer.m_fee_filter_sent = filterToSend;
5836 (currentFilter < 3 * peer.m_fee_filter_sent / 4 || currentFilter > 4 * peer.m_fee_filter_sent / 3)) {
5841bool PeerManagerImpl::RejectIncomingTxs(
const CNode& peer)
const
5854 const size_t nAvail{vRecv.
size()};
5855 bool bPingFinished =
false;
5856 std::string sProblem;
5858 if (nAvail >=
sizeof(
nonce)) {
5862 if (peer.m_ping_nonce_sent != 0) {
5863 if (
nonce == peer.m_ping_nonce_sent) {
5865 bPingFinished =
true;
5866 const auto ping_time = ping_end - peer.m_ping_start.load();
5867 if (ping_time.count() >= 0) {
5871 m_tx_for_private_broadcast.NodeConfirmedReception(pfrom.
GetId());
5878 sProblem =
"Timing mishap";
5882 sProblem =
"Nonce mismatch";
5885 bPingFinished =
true;
5886 sProblem =
"Nonce zero";
5890 sProblem =
"Unsolicited pong without ping";
5894 bPingFinished =
true;
5895 sProblem =
"Short payload";
5898 if (!(sProblem.empty())) {
5902 peer.m_ping_nonce_sent,
5906 if (bPingFinished) {
5907 peer.m_ping_nonce_sent = 0;
5911bool PeerManagerImpl::SetupAddressRelay(
const CNode&
node, Peer& peer)
5916 if (
node.IsBlockOnlyConn())
return false;
5921 if (
node.IsFeelerConn())
return false;
5923 if (!peer.m_addr_relay_enabled.exchange(
true)) {
5927 peer.m_addr_known = std::make_unique<CRollingBloomFilter>(5000, 0.001);
5933void PeerManagerImpl::ProcessAddrs(std::string_view msg_type,
CNode& pfrom, Peer& peer, std::vector<CAddress>&& vAddr,
const std::atomic<bool>& interruptMsgProc)
5938 if (!SetupAddressRelay(pfrom, peer)) {
5945 Misbehaving(peer,
strprintf(
"%s message size = %u", msg_type, vAddr.size()));
5950 std::vector<CAddress> vAddrOk;
5956 const auto time_diff{current_time - peer.m_addr_token_timestamp};
5960 peer.m_addr_token_timestamp = current_time;
5963 uint64_t num_proc = 0;
5964 uint64_t num_rate_limit = 0;
5965 std::shuffle(vAddr.begin(), vAddr.end(), m_rng);
5968 if (interruptMsgProc)
5972 if (peer.m_addr_token_bucket < 1.0) {
5978 peer.m_addr_token_bucket -= 1.0;
5987 addr.
nTime = std::chrono::time_point_cast<std::chrono::seconds>(current_time - 5 * 24h);
5989 AddAddressKnown(peer, addr);
5996 if (addr.
nTime > current_time - 10min && !peer.m_getaddr_sent && vAddr.size() <= 10 && addr.
IsRoutable()) {
5998 RelayAddress(pfrom.
GetId(), addr, reachable);
6002 vAddrOk.push_back(addr);
6005 peer.m_addr_processed += num_proc;
6006 peer.m_addr_rate_limited += num_rate_limit;
6007 LogDebug(
BCLog::NET,
"Received addr: %u addresses (%u processed, %u rate-limited) from peer=%d\n",
6008 vAddr.size(), num_proc, num_rate_limit, pfrom.
GetId());
6010 m_addrman.
Add(vAddrOk, pfrom.
addr, 2h);
6011 if (vAddr.size() < 1000) peer.m_getaddr_sent =
false;
6020bool PeerManagerImpl::SendMessages(
CNode&
node)
6025 PeerRef maybe_peer{GetPeerRef(
node.GetId())};
6026 if (!maybe_peer)
return false;
6027 Peer& peer{*maybe_peer};
6032 if (MaybeDiscourageAndDisconnect(
node, peer))
return true;
6035 if (!
node.IsInboundConn() && !peer.m_outbound_version_message_sent) {
6036 PushNodeVersion(
node, peer);
6037 peer.m_outbound_version_message_sent =
true;
6041 if (!
node.fSuccessfullyConnected ||
node.fDisconnect)
6045 const auto current_time{GetTime<std::chrono::microseconds>()};
6050 if (
node.IsPrivateBroadcastConn()) {
6054 node.fDisconnect =
true;
6061 node.fDisconnect =
true;
6065 MaybeSendPing(
node, peer, now);
6068 if (
node.fDisconnect)
return true;
6070 MaybeSendAddr(
node, peer, current_time);
6072 MaybeSendSendHeaders(
node, peer);
6074 ProcessInvBacklog(now);
6079 CNodeState &state = *
State(
node.GetId());
6082 if (m_chainman.m_best_header ==
nullptr) {
6089 bool sync_blocks_and_headers_from_peer =
false;
6090 if (state.fPreferredDownload) {
6091 sync_blocks_and_headers_from_peer =
true;
6092 }
else if (CanServeBlocks(peer) && !
node.IsAddrFetchConn()) {
6102 if (m_num_preferred_download_peers == 0 || mapBlocksInFlight.empty()) {
6103 sync_blocks_and_headers_from_peer =
true;
6109 if ((nSyncStarted == 0 && sync_blocks_and_headers_from_peer) || m_chainman.m_best_header->Time() >
NodeClock::now() - 24h) {
6110 const CBlockIndex* pindexStart = m_chainman.m_best_header;
6118 if (pindexStart->
pprev)
6119 pindexStart = pindexStart->
pprev;
6123 state.fSyncStarted =
true;
6147 LOCK(peer.m_block_inv_mutex);
6148 std::vector<CBlock> vHeaders;
6149 bool fRevertToInv = ((!peer.m_prefers_headers &&
6150 (!state.m_requested_hb_cmpctblocks || peer.m_blocks_for_headers_relay.size() > 1)) ||
6153 ProcessBlockAvailability(
node.GetId());
6155 if (!fRevertToInv) {
6156 bool fFoundStartingHeader =
false;
6160 for (
const uint256& hash : peer.m_blocks_for_headers_relay) {
6165 fRevertToInv =
true;
6168 if (pBestIndex !=
nullptr && pindex->
pprev != pBestIndex) {
6180 fRevertToInv =
true;
6183 pBestIndex = pindex;
6184 if (fFoundStartingHeader) {
6187 }
else if (PeerHasHeader(&state, pindex)) {
6189 }
else if (pindex->
pprev ==
nullptr || PeerHasHeader(&state, pindex->
pprev)) {
6192 fFoundStartingHeader =
true;
6197 fRevertToInv =
true;
6202 if (!fRevertToInv && !vHeaders.empty()) {
6203 if (vHeaders.size() == 1 && state.m_requested_hb_cmpctblocks) {
6207 vHeaders.front().GetHash().ToString(),
node.GetId());
6209 std::optional<CSerializedNetMsg> cached_cmpctblock_msg;
6211 LOCK(m_most_recent_block_mutex);
6212 if (m_most_recent_block_hash == pBestIndex->
GetBlockHash()) {
6216 if (cached_cmpctblock_msg.has_value()) {
6217 PushMessage(
node, std::move(cached_cmpctblock_msg.value()));
6225 state.pindexBestHeaderSent = pBestIndex;
6226 }
else if (peer.m_prefers_headers) {
6227 if (vHeaders.size() > 1) {
6230 vHeaders.front().GetHash().ToString(),
6231 vHeaders.back().GetHash().ToString(),
node.GetId());
6234 vHeaders.front().GetHash().ToString(),
node.GetId());
6237 state.pindexBestHeaderSent = pBestIndex;
6239 fRevertToInv =
true;
6245 if (!peer.m_blocks_for_headers_relay.empty()) {
6246 const uint256& hashToAnnounce = peer.m_blocks_for_headers_relay.back();
6259 if (!PeerHasHeader(&state, pindex)) {
6260 peer.m_blocks_for_inv_relay.push_back(hashToAnnounce);
6266 peer.m_blocks_for_headers_relay.clear();
6272 std::vector<CInv> vInv;
6274 LOCK(peer.m_block_inv_mutex);
6275 vInv.reserve(peer.m_blocks_for_inv_relay.size());
6278 for (
const uint256& hash : peer.m_blocks_for_inv_relay) {
6285 peer.m_blocks_for_inv_relay.clear();
6288 if (
auto tx_relay = peer.GetTxRelay(); tx_relay !=
nullptr) {
6289 LOCK(tx_relay->m_tx_inventory_mutex);
6292 if (tx_relay->m_next_inv_send_time < current_time) {
6293 fSendTrickle =
true;
6294 if (
node.IsInboundConn()) {
6303 LOCK(tx_relay->m_bloom_filter_mutex);
6304 if (!tx_relay->m_relay_txs) tx_relay->m_tx_inventory_to_send.clear();
6308 if (fSendTrickle && tx_relay->m_send_mempool) {
6309 auto vtxinfo = m_mempool.
infoAll();
6314 tx_relay->m_send_mempool =
false;
6315 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6318 tx_relay->m_tx_inventory_to_send.clear();
6320 LOCK(tx_relay->m_bloom_filter_mutex);
6322 for (
const auto& txinfo : vtxinfo) {
6323 const Txid& txid{txinfo.tx->GetHash()};
6324 const Wtxid& wtxid{txinfo.tx->GetWitnessHash()};
6325 const auto inv = peer.m_wtxid_relay ?
6330 if (txinfo.fee < filterrate.GetFee(txinfo.vsize)) {
6333 if (tx_relay->m_bloom_filter) {
6334 if (!tx_relay->m_bloom_filter->IsRelevantAndUpdate(*txinfo.tx))
continue;
6336 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6337 vInv.push_back(inv);
6349 const CFeeRate filterrate{tx_relay->m_fee_filter_received.load()};
6352 auto& invs = tx_relay->m_tx_inventory_to_send;
6353 std::vector<CTransactionRef> res;
6355 if (invs.size() == 0)
return res;
6358 if (invs.capacity() > 2 * invs.size()) invs.shrink_to_fit();
6362 res.reserve(txiters.size());
6363 for (
auto txiter : txiters) {
6364 if (txiter->GetFee() < filterrate.GetFee(txiter->GetTxSize())) {
6367 res.push_back(txiter->GetSharedTx());
6370 tx_relay->m_last_inv_sequence = m_mempool.
GetSequence();
6374 LOCK(tx_relay->m_bloom_filter_mutex);
6375 vInv.reserve(std::min<size_t>(
MAX_INV_SZ, vInv.size() + inv_tx.size()));
6376 for (
auto& tx : inv_tx) {
6380 const auto inv = peer.m_wtxid_relay ?
6384 if (tx_relay->m_tx_inventory_known_filter.contains(inv.
hash)) {
6387 if (tx_relay->m_bloom_filter && !tx_relay->m_bloom_filter->IsRelevantAndUpdate(*tx))
continue;
6389 vInv.push_back(inv);
6394 tx_relay->m_tx_inventory_known_filter.insert(inv.
hash);
6402 auto stalling_timeout = m_block_stalling_timeout.load();
6403 if (state.m_stalling_since.count() && state.m_stalling_since < current_time - stalling_timeout) {
6407 LogInfo(
"Peer is stalling block download, %s",
node.DisconnectMsg());
6408 node.fDisconnect =
true;
6412 if (stalling_timeout != new_timeout && m_block_stalling_timeout.compare_exchange_strong(stalling_timeout, new_timeout)) {
6422 if (state.vBlocksInFlight.size() > 0) {
6423 QueuedBlock &queuedBlock = state.vBlocksInFlight.front();
6424 int nOtherPeersWithValidatedDownloads = m_peers_downloading_from - 1;
6426 LogInfo(
"Timeout downloading block %s, %s", queuedBlock.pindex->GetBlockHash().ToString(),
node.DisconnectMsg());
6427 node.fDisconnect =
true;
6432 if (state.fSyncStarted && peer.m_headers_sync_timeout < std::chrono::microseconds::max()) {
6434 if (m_chainman.m_best_header->Time() <=
NodeClock::now() - 24h) {
6435 if (current_time > peer.m_headers_sync_timeout && nSyncStarted == 1 && (m_num_preferred_download_peers - state.fPreferredDownload >= 1)) {
6442 LogInfo(
"Timeout downloading headers, %s",
node.DisconnectMsg());
6443 node.fDisconnect =
true;
6446 LogInfo(
"Timeout downloading headers from noban peer, not %s",
node.DisconnectMsg());
6452 state.fSyncStarted =
false;
6454 peer.m_headers_sync_timeout = 0us;
6460 peer.m_headers_sync_timeout = std::chrono::microseconds::max();
6466 ConsiderEviction(
node, peer, GetTime<std::chrono::seconds>());
6471 std::vector<CInv> vGetData;
6473 std::vector<const CBlockIndex*> vToDownload;
6475 auto get_inflight_budget = [&state]() {
6482 FindNextBlocksToDownload(peer, get_inflight_budget(), vToDownload, staller);
6483 auto historical_blocks{m_chainman.GetHistoricalBlockRange()};
6484 if (historical_blocks && !IsLimitedPeer(peer)) {
6488 TryDownloadingHistoricalBlocks(
6490 get_inflight_budget(),
6491 vToDownload, from_tip, historical_blocks->second);
6494 uint32_t nFetchFlags = GetFetchFlags(peer);
6496 BlockRequested(
node.GetId(), *pindex);
6500 if (state.vBlocksInFlight.empty() && staller != -1) {
6501 if (
State(staller)->m_stalling_since == 0us) {
6502 State(staller)->m_stalling_since = current_time;
6512 LOCK(m_tx_download_mutex);
6513 for (
const GenTxid& gtxid : m_txdownloadman.GetRequestsToSend(
node.GetId(), current_time)) {
6522 if (!vGetData.empty())
6525 MaybeSendFeefilter(
node, peer, current_time);
static constexpr CAmount MAX_MONEY
No amount larger than this (in satoshi) is valid.
bool MoneyRange(const CAmount &nValue)
int64_t CAmount
Amount in satoshis (Can be negative)
enum ReadStatus_t ReadStatus
const std::string & BlockFilterTypeName(BlockFilterType filter_type)
Get the human-readable name for a filter type.
BlockFilterIndex * GetBlockFilterIndex(BlockFilterType filter_type)
Get a block filter index by type.
static constexpr int CFCHECKPT_INTERVAL
Interval between compact filter checkpoints.
CBlockLocator GetLocator(const CBlockIndex *index)
Get a locator for a block index entry.
int64_t GetBlockProofEquivalentTime(const CBlockIndex &to, const CBlockIndex &from, const CBlockIndex &tip, const Consensus::Params ¶ms)
Return the time it would take to redo the work difference between from and to, assuming the current h...
const CBlockIndex * LastCommonAncestor(const CBlockIndex *pa, const CBlockIndex *pb)
Find the last common ancestor two blocks have.
@ BLOCK_VALID_CHAIN
Outputs do not overspend inputs, no double spends, coinbase output ok, no immature coinbase spends,...
@ BLOCK_VALID_TRANSACTIONS
Only first tx is coinbase, 2 <= coinbase input script length <= 100, transactions valid,...
@ BLOCK_VALID_SCRIPTS
Scripts & signatures ok.
@ BLOCK_VALID_TREE
All parent headers found, difficulty matches, timestamp >= median previous.
@ BLOCK_HAVE_DATA
full block available in blk*.dat
arith_uint256 GetBlockProof(const CBlockIndex &block)
Compute how much work a block index entry corresponds to.
#define Assert(val)
Identity function.
#define Assume(val)
Assume is the identity function.
Stochastic address manager.
void Connected(const CService &addr, NodeSeconds time=Now< NodeSeconds >())
We have successfully connected to this peer.
bool Good(const CService &addr, NodeSeconds time=Now< NodeSeconds >())
Mark an address record as accessible and attempt to move it to addrman's tried table.
bool Add(const std::vector< CAddress > &vAddr, const CNetAddr &source, std::chrono::seconds time_penalty=0s)
Attempt to add one or more addresses to addrman's new table.
void SetServices(const CService &addr, ServiceFlags nServices)
Update an entry's service bits.
bool IsBanned(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
Return whether net_addr is banned.
bool IsDiscouraged(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
Return whether net_addr is discouraged.
void Discourage(const CNetAddr &net_addr) EXCLUSIVE_LOCKS_REQUIRED(!m_banned_mutex)
BlockFilterIndex is used to store and retrieve block filters, hashes, and headers for a range of bloc...
bool LookupFilterRange(int start_height, const CBlockIndex *stop_index, std::vector< BlockFilter > &filters_out) const
Get a range of filters between two heights on a chain.
bool LookupFilterHashRange(int start_height, const CBlockIndex *stop_index, std::vector< uint256 > &hashes_out) const
Get a range of filter hashes between two heights on a chain.
bool LookupFilterHeader(const CBlockIndex *block_index, uint256 &header_out) EXCLUSIVE_LOCKS_REQUIRED(!m_cs_headers_cache)
Get a single filter header by block.
std::vector< CTransactionRef > txn
std::vector< uint16_t > indexes
A CService with information about it as peer.
ServiceFlags nServices
Serialized as uint64_t in V1, and as CompactSize in V2.
static constexpr SerParams V1_NETWORK
NodeSeconds nTime
Always included in serialization. The behavior is unspecified if the value is not representable as ui...
static constexpr SerParams V2_NETWORK
size_t BlockTxCount() const
std::vector< CTransactionRef > vtx
The block chain is a tree shaped structure starting with the genesis block at the root,...
bool IsValid(enum BlockStatus nUpTo) const EXCLUSIVE_LOCKS_REQUIRED(
Check whether this block index entry is valid up to the passed validity level.
CBlockIndex * pprev
pointer to the index of the predecessor of this block
CBlockHeader GetBlockHeader() const
arith_uint256 nChainWork
(memory only) Total amount of work (expected number of hashes) in the chain up to and including this ...
bool HaveNumChainTxs() const
Check whether this block and all previous blocks back to the genesis block or an assumeutxo snapshot ...
uint256 GetBlockHash() const
int64_t GetBlockTime() const
unsigned int nTx
Number of transactions in this block.
CBlockIndex * GetAncestor(int height)
Efficiently find an ancestor of this block.
int nHeight
height of the entry in the chain. The genesis block has height 0
FlatFilePos GetBlockPos() const EXCLUSIVE_LOCKS_REQUIRED(
BloomFilter is a probabilistic filter which SPV clients provide so that we can filter the transaction...
bool IsWithinSizeConstraints() const
True if the size is <= MAX_BLOOM_FILTER_SIZE and the number of hash functions is <= MAX_HASH_FUNCS (c...
An in-memory indexed chain of blocks.
bool Contains(const CBlockIndex &index) const
Efficiently check whether a block is present in this chain.
CBlockIndex * Tip() const
Returns the index entry for the tip of this chain, or nullptr if none.
CBlockIndex * Next(const CBlockIndex &index) const
Find the successor of a block in this chain, or nullptr if the given index is not found or is the tip...
int Height() const
Return the maximal height in the chain.
CChainParams defines various tweakable parameters of a given instance of the Bitcoin system.
const HeadersSyncParams & HeadersSync() const
const Consensus::Params & GetConsensus() const
void NumToOpenAdd(size_t n)
Increment the number of new connections of type ConnectionType::PRIVATE_BROADCAST to be opened by CCo...
size_t NumToOpenSub(size_t n)
Decrement the number of new connections of type ConnectionType::PRIVATE_BROADCAST to be opened by CCo...
bool GetNetworkActive() const
bool GetTryNewOutboundPeer() const
class CConnman::PrivateBroadcast m_private_broadcast
bool ShouldRunInactivityChecks(const CNode &node, NodeClock::time_point now) const
Return true if we should disconnect the peer for failing an inactivity check.
std::vector< CAddress > GetAddresses(CNode &requestor, size_t max_addresses, size_t max_pct)
Return addresses from the per-requestor cache.
void SetTryNewOutboundPeer(bool flag)
void WakeMessageHandler() EXCLUSIVE_LOCKS_REQUIRED(!mutexMsgProc)
bool OutboundTargetReached(bool historicalBlockServingLimit) const EXCLUSIVE_LOCKS_REQUIRED(!m_total_bytes_sent_mutex)
check if the outbound target is reached if param historicalBlockServingLimit is set true,...
void StartExtraBlockRelayPeers()
void ForEachNode(const NodeFn &func) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
CSipHasher GetDeterministicRandomizer(uint64_t id) const
Get a unique deterministic randomizer.
bool EvictTxPeerIfFull(std::optional< NodeId > protect_peer=std::nullopt) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
If we are at capacity for inbound tx-relay peers, attempt to evict one.
uint32_t GetMappedAS(const CNetAddr &addr) const
bool ForNode(NodeId id, std::function< bool(CNode *pnode)> func) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
std::vector< CAddress > GetAddressesUnsafe(size_t max_addresses, size_t max_pct, std::optional< Network > network, bool filtered=true) const
Return randomly selected addresses.
int GetExtraBlockRelayCount() const EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
bool DisconnectNode(std::string_view node) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
Mutex & GetNodesMutex() const LOCK_RETURNED(m_nodes_mutex)
int GetExtraFullOutboundCount() const EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
bool GetUseAddrmanOutgoing() const
bool CheckIncomingNonce(uint64_t nonce) EXCLUSIVE_LOCKS_REQUIRED(!m_nodes_mutex)
Fee rate in satoshis per virtualbyte: CAmount / vB the feerate is represented internally as FeeFrac.
CAmount GetFeePerK() const
Return the fee in satoshis for a vsize of 1000 vbytes.
bool IsMsgCmpctBlk() const
std::string ToString() const
bool IsMsgFilteredBlk() const
bool IsMsgWitnessBlk() const
Used to relay blocks as header + vector<merkle branch> to filtered nodes.
std::vector< std::pair< unsigned int, Txid > > vMatchedTxn
Public only for unit testing and relay testing (not relayed).
bool IsRelayable() const
Whether this address should be relayed to other peers even if we can't reach it ourselves.
static constexpr SerParams V1
bool IsAddrV1Compatible() const
Check if the current object can be serialized in pre-ADDRv2/BIP155 format.
Transport protocol agnostic message container.
Information about a peer.
bool IsFeelerConn() const
bool ExpectServicesFromConn() const
std::atomic< int > nVersion
std::atomic_bool m_has_all_wanted_services
Whether this peer provides all services that we want.
bool IsInboundConn() const
bool HasPermission(NetPermissionFlags permission) const
std::string LogPeer() const
Helper function to log the peer id, optionally including IP address.
bool IsOutboundOrBlockRelayConn() const
bool IsManualConn() const
std::string ConnectionTypeAsString() const
void SetCommonVersion(int greatest_common_version)
std::atomic< bool > m_bip152_highbandwidth_to
std::atomic_bool m_relays_txs
Whether we should relay transactions to this peer.
std::atomic< bool > m_bip152_highbandwidth_from
std::atomic_bool fSuccessfullyConnected
fSuccessfullyConnected is set to true on receiving VERACK from the peer.
bool IsAddrFetchConn() const
uint64_t GetLocalNonce() const
void SetAddrLocal(const CService &addrLocalIn) EXCLUSIVE_LOCKS_REQUIRED(!m_addr_local_mutex)
May not be called more than once.
const NodeClock::time_point m_connected
Unix epoch time at peer connection.
bool IsBlockOnlyConn() const
int GetCommonVersion() const
bool IsFullOutboundConn() const
std::atomic_bool fPauseSend
std::string DisconnectMsg() const
Helper function to log disconnects.
void PongReceived(NodeClock::duration ping_time)
A ping-pong round trip has completed successfully. Update latest and minimum ping durations.
std::atomic_bool m_bloom_filter_loaded
Whether this peer has loaded a bloom filter.
bool IsPrivateBroadcastConn() const
const std::unique_ptr< Transport > m_transport
Transport serializer/deserializer.
const bool m_inbound_onion
Whether this peer is an inbound onion, i.e. connected via our Tor onion service.
std::atomic< std::chrono::seconds > m_last_block_time
UNIX epoch time of the last block received from this peer that we had not yet seen (e....
int AdvertisedVersion() const
Protocol version advertised in our VERSION message.
std::atomic_bool fDisconnect
std::atomic< std::chrono::seconds > m_last_tx_time
UNIX epoch time of the last transaction received from this peer that we had not yet seen (e....
RollingBloomFilter is a probabilistic "keep track of most recently inserted" set.
Simple class for background tasks that should be run periodically or once "after a while".
void scheduleEvery(Function f, std::chrono::milliseconds delta) EXCLUSIVE_LOCKS_REQUIRED(!newTaskMutex)
Repeat f until the scheduler is stopped.
void scheduleFromNow(Function f, std::chrono::milliseconds delta) EXCLUSIVE_LOCKS_REQUIRED(!newTaskMutex)
Call f once after the delta has passed.
A combination of a network address (CNetAddr) and a (TCP) port.
std::string ToStringAddrPort() const
std::vector< unsigned char > GetKey() const
General SipHash-2-4 implementation.
uint64_t Finalize() const
Compute the 64-bit SipHash-2-4 of the data written so far.
CSipHasher & Write(uint64_t data)
Hash a 64-bit integer worth of data.
CTxMemPool stores valid-according-to-the-current-best-chain transactions that may be included in the ...
TxMempoolInfo info_for_relay(const T &id, uint64_t last_sequence) const
Returns info for a transaction if its entry_sequence < last_sequence.
CFeeRate GetMinFee(size_t sizelimit) const
CTransactionRef get(const Txid &hash) const
Return a mempool transaction with a given hash.
size_t DynamicMemoryUsage() const
std::vector< TxMempoolInfo > infoAll() const
bool exists(const Txid &txid) const
uint64_t GetSequence() const EXCLUSIVE_LOCKS_REQUIRED(cs)
std::set< Txid > GetUnbroadcastTxs() const
Returns transactions in unbroadcast set.
unsigned long size() const
void RemoveUnbroadcastTx(const Txid &txid, bool unchecked=false)
Removes a transaction from the unbroadcast set.
std::vector< txiter > ExtractBestByMiningScoreWithTopology(std::vector< Wtxid > &wtxids, size_t n_to_sort) const EXCLUSIVE_LOCKS_REQUIRED(cs)
Look up wtxids in the mempool and (partially) sort by mining score.
virtual void NewPoWValidBlock(const CBlockIndex *pindex, const std::shared_ptr< const CBlock > &block)
Notifies listeners that a block which builds directly on our current tip has been received and connec...
virtual void UpdatedBlockTip(const CBlockIndex *pindexNew, const CBlockIndex *pindexFork, bool fInitialDownload)
Notifies listeners when the block chain tip advances.
virtual void BlockChecked(const std::shared_ptr< const CBlock > &, const BlockValidationState &)
Notifies listeners of a block validation result.
virtual void ActiveTipChange(const CBlockIndex &new_tip, bool is_ibd)
Notifies listeners any time the block chain tip changes, synchronously.
virtual void BlockDisconnected(const std::shared_ptr< const CBlock > &block, const CBlockIndex *pindex)
Notifies listeners of a block being disconnected Provides the block that was disconnected.
virtual void BlockConnected(const kernel::ChainstateRole &role, const std::shared_ptr< const CBlock > &block, const CBlockIndex *pindex)
Notifies listeners of a block being connected.
void ClearBlockIndexCandidates() EXCLUSIVE_LOCKS_REQUIRED(void PopulateBlockIndexCandidates() EXCLUSIVE_LOCKS_REQUIRED(const CBlockIndex * FindForkInGlobalIndex(const CBlockLocator &locator) const EXCLUSIVE_LOCKS_REQUIRED(cs_main)
Populate the candidate set by calling TryAddBlockIndexCandidate on all valid block indices.
Interface for managing multiple Chainstate objects, where each chainstate is associated with chainsta...
bool IsInitialBlockDownload() const noexcept
Check whether we are doing an initial block download (synchronizing from disk or network)
MempoolAcceptResult ProcessTransaction(const CTransactionRef &tx, bool test_accept=false) EXCLUSIVE_LOCKS_REQUIRED(cs_main)
Try to add a transaction to the memory pool.
RecursiveMutex & GetMutex() const LOCK_RETURNED(
Alias for cs_main.
CBlockIndex * ActiveTip() const EXCLUSIVE_LOCKS_REQUIRED(GetMutex())
Chainstate & ActiveChainstate() const
Alternatives to CurrentChainstate() used by older code to query latest chainstate information without...
SnapshotCompletionResult MaybeValidateSnapshot(Chainstate &validated_cs, Chainstate &unvalidated_cs) EXCLUSIVE_LOCKS_REQUIRED(Chainstate & CurrentChainstate() const EXCLUSIVE_LOCKS_REQUIRED(GetMutex())
Try to validate an assumeutxo snapshot by using a validated historical chainstate targeted at the sna...
bool ProcessNewBlock(const std::shared_ptr< const CBlock > &block, bool force_processing, bool min_pow_checked, bool *new_block) LOCKS_EXCLUDED(cs_main)
Process an incoming block.
bool ProcessNewBlockHeaders(std::span< const CBlockHeader > headers, bool min_pow_checked, BlockValidationState &state, const CBlockIndex **ppindex=nullptr) LOCKS_EXCLUDED(cs_main)
Process incoming block headers.
const arith_uint256 & MinimumChainWork() const
CChain & ActiveChain() const EXCLUSIVE_LOCKS_REQUIRED(GetMutex())
void ReportHeadersPresync(int64_t height, int64_t timestamp)
This is used by net_processing to report pre-synchronization progress of headers, as headers are not ...
node::BlockManager m_blockman
A single BlockManager instance is shared across each constructed chainstate to avoid duplicating bloc...
Double ended buffer combining vector and stream-like interfaces.
void ignore(size_t num_ignore)
uint64_t rand64() noexcept
Generate a random 64-bit integer.
const uint256 & ToUint256() const LIFETIMEBOUND
static Mutex g_msgproc_mutex
Mutex for anything that is only accessed via the msg processing thread.
virtual void FinalizeNode(const CNode &node)=0
Handle removal of a peer (clear state)
virtual bool ProcessMessages(CNode &node, std::atomic< bool > &interrupt) EXCLUSIVE_LOCKS_REQUIRED(g_msgproc_mutex)=0
Process protocol messages received from a given node.
virtual bool HasAllDesirableServiceFlags(ServiceFlags services) const =0
Callback to determine whether the given set of service flags are sufficient for a peer to be "relevan...
virtual bool SendMessages(CNode &node) EXCLUSIVE_LOCKS_REQUIRED(g_msgproc_mutex)=0
Send queued protocol messages to a given node.
virtual void InitializeNode(const CNode &node, ServiceFlags our_services)=0
Initialize a peer (setup state)
static bool HasFlag(NetPermissionFlags flags, NetPermissionFlags f)
ReadStatus FillBlock(CBlock &block, const std::vector< CTransactionRef > &vtx_missing, bool segwit_active)
bool IsTxAvailable(size_t index) const
ReadStatus InitData(const CBlockHeaderAndShortTxIDs &cmpctblock, const std::vector< std::pair< Wtxid, CTransactionRef > > &extra_txn)
virtual util::Expected< void, std::string > FetchBlock(NodeId peer_id, const CBlockIndex &block_index)=0
Attempt to manually fetch block from a given peer.
virtual ServiceFlags GetDesirableServiceFlags(ServiceFlags services) const =0
Gets the set of service flags which are "desirable" for a given peer.
virtual void StartScheduledTasks(CScheduler &scheduler)=0
Begin running background tasks, should only be called once.
virtual std::vector< node::TxOrphanage::OrphanInfo > GetOrphanTransactions()=0
static std::unique_ptr< PeerManager > make(CConnman &connman, AddrMan &addrman, BanMan *banman, ChainstateManager &chainman, CTxMemPool &pool, node::Warnings &warnings, Options opts)
virtual void UnitTestMisbehaving(NodeId peer_id)=0
virtual bool GetNodeStateStats(NodeId nodeid, CNodeStateStats &stats) const =0
Get statistics from node state.
virtual void UpdateLastBlockAnnounceTime(NodeId node, int64_t time_in_seconds)=0
This function is used for testing the stale tip eviction logic, see denialofservice_tests....
virtual void CheckForStaleTipAndEvictPeers()=0
Evict extra outbound peers.
Store a list of transactions to be broadcast privately.
@ QueueFull
Rejected: the queue is already at MAX_TRANSACTIONS.
@ AlreadyPresent
The transaction was already present; no change.
@ Added
The transaction was newly added.
static constexpr size_t MAX_TRANSACTIONS
Maximum number of transactions tracked simultaneously.
I randrange(I range) noexcept
Generate a random integer in the range [0..range), with range > 0.
bool Contains(Network net) const EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
std::string GetDebugMessage() const
std::string ToString() const
256-bit unsigned big integer.
constexpr bool IsNull() const
std::string ToString() const
CBlockIndex * LookupBlockIndex(const uint256 &hash) EXCLUSIVE_LOCKS_REQUIRED(cs_main)
CBlockFileInfo *GetBlockFileInfo(size_t n) EXCLUSIVE_LOCKS_REQUIRED(bool WriteBlockUndo(const CBlockUndo &blockundo, BlockValidationState &state, CBlockIndex &block) EXCLUSIVE_LOCKS_REQUIRED(FlatFilePos WriteBlock(const CBlock &block, int nHeight) EXCLUSIVE_LOCKS_REQUIRED(void UpdateBlockInfo(const CBlock &block, unsigned int nHeight, const FlatFilePos &pos) EXCLUSIVE_LOCKS_REQUIRED(bool IsPruneMode() const
Get block file info entry for one block file.
bool LoadingBlocks() const
ReadRawBlockResult ReadRawBlock(const FlatFilePos &pos, std::optional< std::pair< size_t, size_t > > block_part=std::nullopt) const
bool ReadBlock(CBlock &block, const FlatFilePos &pos, const std::optional< uint256 > &expected_hash) const
Functions for disk access for blocks.
Class responsible for deciding what transactions to request and, once downloaded, whether and how to ...
Manages warning messages within a node.
transaction_identifier represents the two canonical transaction identifier types (txid,...
std::string ToString() const
const uint256 & ToUint256() const LIFETIMEBOUND
The util::Expected class provides a standard way for low-level functions to return either error value...
A token bucket rate limiter.
bool decrement(double n=1.0, double floor=0.0)
Consume n tokens.
void increment(const time_point &now)
Refill tokens based on elapsed time since last call.
double value() const
Current token balance.
The util::Unexpected class represents an unexpected value stored in util::Expected.
std::string TransportTypeAsString(TransportProtocolType transport_type)
Convert TransportProtocolType enum to a string value.
@ BLOCK_HEADER_LOW_WORK
the block header may be on a too-little-work chain
@ BLOCK_INVALID_HEADER
invalid proof of work or time too old
@ BLOCK_CACHED_INVALID
this block was cached as being invalid and we didn't store the reason why
@ BLOCK_CONSENSUS
invalid by consensus rules (excluding any below reasons)
@ BLOCK_MISSING_PREV
We don't have the previous block the checked one is built on.
@ BLOCK_INVALID_PREV
A block this one builds on is invalid.
@ BLOCK_MUTATED
the block's data didn't match the data committed to by the PoW
@ BLOCK_TIME_FUTURE
block timestamp was > 2 hours in the future (or our clock is bad)
@ BLOCK_RESULT_UNSET
initial value. Block has not yet been rejected
@ TX_MISSING_INPUTS
transaction was missing some of its inputs
@ TX_UNKNOWN
transaction was not validated because package failed
@ TX_NO_MEMPOOL
this node does not have a mempool so can't validate the transaction
@ TX_RESULT_UNSET
initial value. Tx has not yet been rejected
static size_t RecursiveDynamicUsage(const CScript &script)
RecursiveMutex cs_main
Mutex to guard access to validation specific variables, such as reading or changing the chainstate.
bool DeploymentActiveAfter(const CBlockIndex *pindexPrev, const Consensus::Params ¶ms, Consensus::BuriedDeployment dep, VersionBitsCache &versionbitscache)
Determine if a deployment is active for the next block.
bool DeploymentActiveAt(const CBlockIndex &index, const Consensus::Params ¶ms, Consensus::BuriedDeployment dep, VersionBitsCache &versionbitscache)
Determine if a deployment is active for this block.
is a home for simple enum and struct type definitions that can be used internally by functions in the...
#define LogDebug(category,...)
CSerializedNetMsg Make(std::string msg_type, Args &&... args)
constexpr const char * FILTERCLEAR
The filterclear message tells the receiving peer to remove a previously-set bloom filter.
constexpr const char * FEEFILTER
The feefilter message tells the receiving peer not to inv us any txs which do not meet the specified ...
constexpr const char * SENDHEADERS
Indicates that a node prefers to receive new block announcements via a "headers" message rather than ...
constexpr const char * GETBLOCKS
The getblocks message requests an inv message that provides block header hashes starting from a parti...
constexpr const char * HEADERS
The headers message sends one or more block headers to a node which previously requested certain head...
constexpr const char * ADDR
The addr (IP address) message relays connection information for peers on the network.
constexpr const char * GETBLOCKTXN
Contains a BlockTransactionsRequest Peer should respond with "blocktxn" message.
constexpr const char * CMPCTBLOCK
Contains a CBlockHeaderAndShortTxIDs object - providing a header and list of "short txids".
constexpr const char * CFCHECKPT
cfcheckpt is a response to a getcfcheckpt request containing a vector of evenly spaced filter headers...
constexpr const char * SENDADDRV2
The sendaddrv2 message signals support for receiving ADDRV2 messages (BIP155).
constexpr const char * GETADDR
The getaddr message requests an addr message from the receiving node, preferably one with lots of IP ...
constexpr const char * GETCFILTERS
getcfilters requests compact filters for a range of blocks.
constexpr const char * PONG
The pong message replies to a ping message, proving to the pinging node that the ponging node is stil...
constexpr const char * BLOCKTXN
Contains a BlockTransactions.
constexpr const char * CFHEADERS
cfheaders is a response to a getcfheaders request containing a filter header and a vector of filter h...
constexpr const char * PING
The ping message is sent periodically to help confirm that the receiving peer is still connected.
constexpr const char * FILTERLOAD
The filterload message tells the receiving peer to filter all relayed transactions and requested merk...
constexpr const char * SENDTXRCNCL
Contains a 4-byte version number and an 8-byte salt.
constexpr const char * ADDRV2
The addrv2 message relays connection information for peers on the network just like the addr message,...
constexpr const char * VERACK
The verack message acknowledges a previously-received version message, informing the connecting node ...
constexpr const char * GETHEADERS
The getheaders message requests a headers message that provides block headers starting from a particu...
constexpr const char * FILTERADD
The filteradd message tells the receiving peer to add a single element to a previously-set bloom filt...
constexpr const char * CFILTER
cfilter is a response to a getcfilters request containing a single compact filter.
constexpr const char * FEATURE
BIP 434 Peer feature negotiation.
constexpr const char * GETDATA
The getdata message requests one or more data objects from another node.
constexpr const char * SENDCMPCT
Contains a 1-byte bool and 8-byte LE version number.
constexpr const char * GETCFCHECKPT
getcfcheckpt requests evenly spaced compact filter headers, enabling parallelized download and valida...
constexpr const char * INV
The inv message (inventory message) transmits one or more inventories of objects known to the transmi...
constexpr const char * TX
The tx message transmits a single transaction.
constexpr const char * MEMPOOL
The mempool message requests the TXIDs of transactions that the receiving node has verified as valid ...
constexpr const char * NOTFOUND
The notfound message is a reply to a getdata message which requested an object the receiving node doe...
constexpr const char * MERKLEBLOCK
The merkleblock message is a reply to a getdata message which requested a block using the inventory t...
constexpr const char * WTXIDRELAY
Indicates that a node prefers to relay transactions via wtxid, rather than txid.
constexpr const char * BLOCK
The block message transmits a single serialized block.
constexpr const char * GETCFHEADERS
getcfheaders requests a compact filter header and the filter hashes for a range of blocks,...
constexpr const char * VERSION
The version message provides information about the transmitting node to the receiving node at the beg...
static constexpr int32_t MAX_PEER_TX_ANNOUNCEMENTS
Maximum number of transactions to consider for requesting, per peer.
""_hex is a compile-time user-defined literal returning a std::array<std::byte>, equivalent to ParseH...
bool ShouldDebugLog(Category category)
Return whether messages with specified category should be debug logged.
std::string ToString(const T &t)
Locale-independent version of std::to_string.
std::string strSubVersion
Subversion as sent to the P2P network in version messages.
std::optional< CService > GetLocalAddrForPeer(CNode &node)
Returns a local address that we should advertise to this peer.
std::function< void(const CAddress &addr, const std::string &msg_type, std::span< const unsigned char > data, bool is_incoming)> CaptureMessage
Defaults to CaptureMessageToFile(), but can be overridden by unit tests.
bool SeenLocal(const CService &addr)
vote for a local address
static const unsigned int MAX_SUBVERSION_LENGTH
Maximum length of the user agent string in version message.
static constexpr std::chrono::minutes TIMEOUT_INTERVAL
Time after which to disconnect, after waiting for a ping response (or inactivity).
static constexpr auto HEADERS_RESPONSE_TIME
How long to wait for a peer to respond to a getheaders request.
static constexpr size_t MAX_ADDR_TO_SEND
The maximum number of address records permitted in an ADDR message.
static constexpr auto INVENTORY_BUCKET_CHECK_DELAY
Delay between checking inventory bucket and backlog.
static constexpr size_t MAX_ADDR_PROCESSING_TOKEN_BUCKET
The soft limit of the address processing token bucket (the regular MAX_ADDR_RATE_PER_SECOND based inc...
static constexpr auto INVENTORY_BUCKET_BACKLOG_HEARTBEAT
Delay between inventory bucket backlog heartbeat log entries.
TRACEPOINT_SEMAPHORE(net, inbound_message)
static const int MAX_BLOCKS_IN_TRANSIT_PER_PEER
Number of blocks that can be requested at any given time from a single peer.
static constexpr auto BLOCK_STALLING_TIMEOUT_DEFAULT
Default time during which a peer must stall block download progress before being disconnected.
static constexpr auto AVG_FEEFILTER_BROADCAST_INTERVAL
Average delay between feefilter broadcasts in seconds.
static constexpr auto EXTRA_PEER_CHECK_INTERVAL
How frequently to check for extra outbound peers and disconnect.
static const unsigned int BLOCK_DOWNLOAD_WINDOW
Size of the "block download window": how far ahead of our current height do we fetch?...
static constexpr int STALE_RELAY_AGE_LIMIT
Age after which a stale block will no longer be served if requested as protection against fingerprint...
static constexpr int HISTORICAL_BLOCK_AGE
Age after which a block is considered historical for purposes of rate limiting block relay.
static constexpr auto ROTATE_ADDR_RELAY_DEST_INTERVAL
Delay between rotating the peers we relay a particular address to.
static constexpr auto MINIMUM_CONNECT_TIME
Minimum time an outbound-peer-eviction candidate must be connected for, in order to evict.
static constexpr auto CHAIN_SYNC_TIMEOUT
Timeout for (unprotected) outbound peers to sync to our chainwork.
static constexpr auto OUTBOUND_INVENTORY_BROADCAST_INTERVAL
Average delay between trickled inventory transmissions for outbound peers.
static const unsigned int NODE_NETWORK_LIMITED_MIN_BLOCKS
Minimum blocks required to signal NODE_NETWORK_LIMITED.
static constexpr auto AVG_LOCAL_ADDRESS_BROADCAST_INTERVAL
Average delay between local address broadcasts.
static const int MAX_BLOCKTXN_DEPTH
Maximum depth of blocks we're willing to respond to GETBLOCKTXN requests for.
static constexpr size_t INVENTORY_BUCKET_BACKLOG_CAPACITY
Empty backlog target capacity.
static constexpr int32_t MAX_OUTBOUND_PEERS_TO_PROTECT_FROM_DISCONNECT
Protect at least this many outbound peers from disconnection due to slow/ behind headers chain.
static constexpr auto INBOUND_INVENTORY_BROADCAST_INTERVAL
Average delay between trickled inventory transmissions for inbound peers.
static constexpr size_t NUM_PRIVATE_BROADCAST_PER_TX
For private broadcast, send a transaction to this many peers.
static constexpr auto MAX_FEEFILTER_CHANGE_DELAY
Maximum feefilter broadcast delay after significant change.
static constexpr uint32_t MAX_GETCFILTERS_SIZE
Maximum number of compact filters that may be requested with one getcfilters.
static constexpr double OUTBOUND_INVENTORY_BUCKET_MULTIPLIER
Multiplier for the inventory bucket rate for outbounds.
static constexpr auto HEADERS_DOWNLOAD_TIMEOUT_BASE
Headers download timeout.
static const unsigned int MAX_GETDATA_SZ
Limit to avoid sending big packets.
static constexpr double BLOCK_DOWNLOAD_TIMEOUT_BASE
Block download timeout base, expressed in multiples of the block interval (i.e.
static constexpr auto PRIVATE_BROADCAST_MAX_CONNECTION_LIFETIME
Private broadcast connections must complete within this time.
static constexpr auto STALE_CHECK_INTERVAL
How frequently to check for stale tips.
static constexpr auto AVG_ADDRESS_BROADCAST_INTERVAL
Average delay between peer address broadcasts.
static const unsigned int MAX_LOCATOR_SZ
The maximum number of entries in a locator.
static constexpr double BLOCK_DOWNLOAD_TIMEOUT_PER_PEER
Additional block download timeout per parallel downloading peer (i.e.
static constexpr double MAX_ADDR_RATE_PER_SECOND
The maximum rate of address records we're willing to process on average.
static constexpr auto PING_INTERVAL
Time between pings automatically sent out for latency probing and keepalive.
static constexpr size_t INVENTORY_BUCKET_BACKLOG_HEARTBEAT_MIN
Minimum backlog to trigger heartbeat log entries.
static const int MAX_CMPCTBLOCK_DEPTH
Maximum depth of blocks we're willing to serve as compact blocks to peers when requested.
static const unsigned int MAX_BLOCKS_TO_ANNOUNCE
Maximum number of headers to announce when relaying blocks with headers message.
static const unsigned int NODE_NETWORK_LIMITED_ALLOW_CONN_BLOCKS
Window, in blocks, for connecting to NODE_NETWORK_LIMITED peers.
static constexpr uint32_t MAX_GETCFHEADERS_SIZE
Maximum number of cf hashes that may be requested with one getcfheaders.
static constexpr auto BLOCK_STALLING_TIMEOUT_MAX
Maximum timeout for stalling block download.
static constexpr auto HEADERS_DOWNLOAD_TIMEOUT_PER_HEADER
static constexpr uint64_t RANDOMIZER_ID_ADDRESS_RELAY
SHA256("main address relay")[0:8].
static constexpr size_t MAX_PCT_ADDR_TO_SEND
the maximum percentage of addresses from our addrman to return in response to a getaddr message.
static const unsigned int MAX_INV_SZ
The maximum number of entries in an 'inv' protocol message.
static constexpr uint64_t CMPCTBLOCKS_VERSION
The compactblocks version we support.
static const unsigned int MAX_CMPCTBLOCKS_INFLIGHT_PER_BLOCK
Maximum number of outstanding CMPCTBLOCK requests for the same block.
ReachableNets g_reachable_nets
bool IsProxy(const CNetAddr &addr)
static constexpr unsigned int DEFAULT_MIN_RELAY_TX_FEE
Default for -minrelaytxfee, minimum relay fee for transactions.
static constexpr TransactionSerParams TX_NO_WITNESS
static constexpr TransactionSerParams TX_WITH_WITNESS
std::shared_ptr< const CTransaction > CTransactionRef
GenTxid ToGenTxid(const CInv &inv)
Convert a TX/WITNESS_TX/WTX CInv to a GenTxid.
const uint32_t MSG_WITNESS_FLAG
getdata message type flags
static constexpr size_t MAX_FEATUREDATA_LENGTH
@ MSG_WTX
Defined in BIP 339.
@ MSG_CMPCT_BLOCK
Defined in BIP152.
@ MSG_WITNESS_BLOCK
Defined in BIP144.
ServiceFlags
nServices flags
static constexpr size_t MAX_FEATUREID_LENGTH
static bool MayHaveUsefulAddressDB(ServiceFlags services)
Checks if a peer with the given service flags may be capable of having a robust address-storage DB.
static const int WTXID_RELAY_VERSION
"wtxidrelay" message type for wtxid-based relay starts with this version
static const int FEATURE_VERSION
"feature" message type for feature negotiation starts with this version
static const int SHORT_IDS_BLOCKS_VERSION
short-id-based block download starts with this version
static const int SENDHEADERS_VERSION
"sendheaders" message type and announcing blocks with headers starts with this version
static const int FEEFILTER_VERSION
"feefilter" tells peers to filter invs to you by fee starts with this version
static const int MIN_PEER_PROTO_VERSION
disconnect from peers older than this proto version
static const int INVALID_CB_NO_BAN_VERSION
not banning for invalid compact blocks starts with this version
static const int BIP0031_VERSION
BIP 0031, pong message, is enabled for all versions AFTER this one.
static const unsigned int MAX_SCRIPT_ELEMENT_SIZE
#define LIMITED_VECTOR(obj, n)
#define LIMITED_STRING(obj, n)
uint64_t ReadCompactSize(Stream &is, bool range_check=true)
Decode a CompactSize-encoded variable-length integer.
constexpr auto MakeUCharSpan(const V &v) -> decltype(UCharSpanCast(std::span{v}))
Like the std::span constructor, but for (const) unsigned char member types only.
Describes a place in the block chain to another node such that if the other node doesn't have the sam...
std::vector< uint256 > vHave
NodeClock::duration m_ping_wait
std::vector< int > vHeightInFlight
CAmount m_fee_filter_received
std::chrono::seconds time_offset
bool m_addr_relay_enabled
uint64_t m_addr_rate_limited
uint64_t m_addr_processed
ServiceFlags their_services
Parameters that influence chain consensus.
int64_t nPowTargetSpacing
std::chrono::seconds PowTargetSpacing() const
Validation result for a transaction evaluated by MemPoolAccept (single or package).
const ResultType m_result_type
Result type.
const TxValidationState m_state
Contains information about why the transaction failed.
@ DIFFERENT_WITNESS
Valid, transaction was already in the mempool.
@ INVALID
Fully validated, valid.
const std::list< CTransactionRef > m_replaced_transactions
Mempool transactions replaced by the tx.
Version of the system clock that is mockable in the context of tests (via FakeNodeClock or SetMockTim...
static time_point now() noexcept
Return current system time or mocked time, if set.
std::chrono::time_point< NodeClock > time_point
static constexpr time_point epoch
Validation result for package mempool acceptance.
PackageValidationState m_state
std::map< Wtxid, MempoolAcceptResult > m_tx_results
Map from wtxid to finished MempoolAcceptResults.
std::chrono::seconds median_outbound_time_offset
Information about chainstate that notifications are sent from.
bool historical
Whether this is a historical chainstate downloading old blocks to validate an assumeutxo snapshot,...
CFeeRate min_relay_feerate
A fee rate smaller than this is considered zero fee (for relaying, mining and transaction creation)
std::vector< NodeId > m_senders
std::string ToString() const
#define AssertLockNotHeld(cs)
#define WITH_LOCK(cs, code)
Run code while locking a mutex.
COutPoint ProcessBlock(const NodeContext &node, const std::shared_ptr< CBlock > &block)
Returns the generated coin (or Null if the block was invalid).
#define EXCLUSIVE_LOCKS_REQUIRED(...)
#define LOCKS_EXCLUDED(...)
#define ACQUIRED_BEFORE(...)
#define TRACEPOINT(context,...)
consteval auto _(util::TranslatedLiteral str)
ReconciliationRegisterResult
static constexpr uint32_t TXRECONCILIATION_VERSION
Supported transaction reconciliation protocol version.
std::string SanitizeString(std::string_view str, int rule)
Remove unsafe chars.
int64_t GetTime()
DEPRECATED Use either ClockType::now() or Now<TimePointType>() if a cast is needed.
constexpr int64_t count_seconds(std::chrono::seconds t)
std::chrono::time_point< NodeClock, std::chrono::seconds > NodeSeconds
PackageMempoolAcceptResult ProcessNewPackage(Chainstate &active_chainstate, CTxMemPool &pool, const Package &package, bool test_accept, const std::optional< CFeeRate > &client_maxfeerate)
Validate (and maybe submit) a package to the mempool.
bool IsBlockMutated(const CBlock &block, bool check_witness_root)
Check if a block has been mutated (with respect to its merkle root and witness commitments).
bool HasValidProofOfWork(std::span< const CBlockHeader > headers, const Consensus::Params &consensusParams)
Check that the proof of work on each blockheader matches the value in nBits.
arith_uint256 CalculateClaimedHeadersWork(std::span< const CBlockHeader > headers)
Return the sum of the claimed work on a given set of headers.
static const unsigned int MIN_BLOCKS_TO_KEEP
Block files containing a block-height within MIN_BLOCKS_TO_KEEP of ActiveChain().Tip() will not be pr...
@ UNVALIDATED
Blocks after an assumeutxo snapshot have been validated but the snapshot itself has not been validate...