From 9389a883f185ec8a821a7cf90b8691dbe0c50f63 Mon Sep 17 00:00:00 2001 From: Sami Ahmed Date: Wed, 13 May 2026 00:44:08 -0700 Subject: [PATCH] Wire CSyncManager + LevelDB->RocksDB chain DB migration CSyncManager extracts the headers-first IBD planner from main.cpp into its own translation unit. main.cpp loses ~570 lines of file-scope state and helper functions; the headers handler, block-delivery latency tracking, stall-recovery, and per-peer Tick cadence now route through g_syncManager. MaybeMigrateLevelDbToRocksDb() is now reachable via -migratechaindb / -migratechaindbforce in init.cpp. Reads from /txleveldb and writes byte-for-byte identical records into /rocksdb via a new CRocksTxDB::WriteRawRecordForMigration() shim over WriteRaw. Co-Authored-By: Claude Opus 4.7 (1M context) --- src/CMakeLists.txt | 2 + src/chaindb_migrate.cpp | 198 +++++++++++ src/chaindb_migrate.h | 14 + src/init.cpp | 11 + src/main.cpp | 744 ++-------------------------------------- src/syncmanager.cpp | 664 +++++++++++++++++++++++++++++++++++ src/syncmanager.h | 58 ++++ src/txdb-rocksdb.h | 9 + 8 files changed, 978 insertions(+), 722 deletions(-) create mode 100644 src/chaindb_migrate.cpp create mode 100644 src/chaindb_migrate.h create mode 100644 src/syncmanager.cpp create mode 100644 src/syncmanager.h diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index cb5738f..e433d6b 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -62,6 +62,8 @@ set(CORE_SOURCES pbkdf2.cpp scrypt.cpp smessage.cpp + syncmanager.cpp + chaindb_migrate.cpp tor_embed_hooks.cpp rest.cpp trianglesrpc.cpp diff --git a/src/chaindb_migrate.cpp b/src/chaindb_migrate.cpp new file mode 100644 index 0000000..3088de1 --- /dev/null +++ b/src/chaindb_migrate.cpp @@ -0,0 +1,198 @@ +// Copyright (c) 2026 The Triangles developers. +// Distributed under the MIT/X11 software license, see the accompanying +// file COPYING or http://www.opensource.org/licenses/mit-license.php. + +#include "chaindb_migrate.h" + +#include "txdb-leveldb.h" +#include "txdb-rocksdb.h" +#include "util.h" + +#include +#include +#include + +namespace fs = std::filesystem; + +namespace { + +struct ChainDbStats +{ + int64_t nRecords = 0; + int64_t nUtxos = 0; + int64_t nUtxoValue = 0; + uint256 hashBestChain = 0; + int nDbFormat = 0; +}; + +bool CollectStats(CTxDBBase& db, ChainDbStats& stats, std::string& strError) +{ + stats = ChainDbStats(); + + auto it = db.NewIterator(); + for (it->Seek(std::string()); it->Valid(); it->Next()) + stats.nRecords++; + + int nUtxos = 0; + stats.nUtxoValue = db.SumUtxoValues(nUtxos); + stats.nUtxos = nUtxos; + db.ReadHashBestChain(stats.hashBestChain); + db.ReadDbFormat(stats.nDbFormat); + + if (stats.nRecords <= 0) { + strError = "source chain database contains no records"; + return false; + } + return true; +} + +bool StatsMatch(const ChainDbStats& src, const ChainDbStats& dst, std::string& strError) +{ + if (src.nRecords != dst.nRecords) { + strError = strprintf("record count mismatch after migration: source=%lld rocksdb=%lld", + (long long)src.nRecords, (long long)dst.nRecords); + return false; + } + if (src.nUtxos != dst.nUtxos || src.nUtxoValue != dst.nUtxoValue) { + strError = strprintf("UTXO mismatch after migration: source=(%lld,%lld) rocksdb=(%lld,%lld)", + (long long)src.nUtxos, (long long)src.nUtxoValue, + (long long)dst.nUtxos, (long long)dst.nUtxoValue); + return false; + } + if (src.hashBestChain != dst.hashBestChain) { + strError = strprintf("best-chain hash mismatch after migration: source=%s rocksdb=%s", + src.hashBestChain.ToString().c_str(), + dst.hashBestChain.ToString().c_str()); + return false; + } + if (src.nDbFormat != dst.nDbFormat) { + strError = strprintf("dbformat mismatch after migration: source=%d rocksdb=%d", + src.nDbFormat, dst.nDbFormat); + return false; + } + return true; +} + +} // namespace + +bool MaybeMigrateLevelDbToRocksDb(bool fForce, std::string& strError) +{ + strError.clear(); + + const fs::path dataDir = GetDataDir(); + const fs::path levelPath = dataDir / "txleveldb"; + const fs::path rocksPath = dataDir / "rocksdb"; + const fs::path markerPath = rocksPath / "MIGRATION_INCOMPLETE"; + + if (!fs::exists(levelPath)) + return true; + + if (fs::exists(rocksPath)) { + if (fs::exists(markerPath)) { + printf("ChainDB migration: removing incomplete previous RocksDB migration\n"); + fs::remove_all(rocksPath); + } + else if (!fForce) + return true; + else { + printf("ChainDB migration: removing existing RocksDB directory due to -migratechaindbforce\n"); + fs::remove_all(rocksPath); + } + } + + printf("ChainDB migration: copying LevelDB chain state to RocksDB...\n"); + printf("ChainDB migration: source=%s destination=%s\n", + levelPath.string().c_str(), rocksPath.string().c_str()); + + try { + fs::create_directories(rocksPath); + { + std::ofstream marker(markerPath); + marker << "RocksDB migration in progress. Safe to delete this directory and retry.\n"; + } + + CTxDB source("r"); + CRocksTxDB destination("c+"); + + ChainDbStats srcStats; + if (!CollectStats(source, srcStats, strError)) { + source.Close(); + destination.Close(); + return false; + } + + if (!destination.TxnBegin()) { + strError = "failed to begin RocksDB migration batch"; + source.Close(); + destination.Close(); + return false; + } + + int64_t nCopied = 0; + auto it = source.NewIterator(); + for (it->Seek(std::string()); it->Valid(); it->Next()) + { + if (!destination.WriteRawRecordForMigration(it->KeyStr(), it->ValueStr())) { + destination.TxnAbort(); + strError = "failed to write migrated record to RocksDB"; + source.Close(); + destination.Close(); + return false; + } + + if (++nCopied % 100000 == 0) + { + if (!destination.TxnCommit()) { + strError = "failed to commit RocksDB migration batch"; + source.Close(); + destination.Close(); + return false; + } + printf("ChainDB migration: copied %lld / %lld records\n", + (long long)nCopied, (long long)srcStats.nRecords); + if (!destination.TxnBegin()) { + strError = "failed to begin RocksDB migration batch"; + source.Close(); + destination.Close(); + return false; + } + } + } + + if (!destination.TxnCommit()) { + strError = "failed to commit final RocksDB migration batch"; + source.Close(); + destination.Close(); + return false; + } + + ChainDbStats dstStats; + if (!CollectStats(destination, dstStats, strError)) { + source.Close(); + destination.Close(); + return false; + } + if (!StatsMatch(srcStats, dstStats, strError)) { + source.Close(); + destination.Close(); + return false; + } + + printf("ChainDB migration: verified %lld records, %lld UTXOs, best=%s\n", + (long long)dstStats.nRecords, + (long long)dstStats.nUtxos, + dstStats.hashBestChain.ToString().substr(0,20).c_str()); + + source.Close(); + destination.Close(); + fs::remove(markerPath); + } + catch (std::exception& e) { + strError = e.what(); + return false; + } + + printf("ChainDB migration: complete. Legacy LevelDB was left untouched at %s\n", + levelPath.string().c_str()); + return true; +} diff --git a/src/chaindb_migrate.h b/src/chaindb_migrate.h new file mode 100644 index 0000000..5862c4d --- /dev/null +++ b/src/chaindb_migrate.h @@ -0,0 +1,14 @@ +// Copyright (c) 2026 The Triangles developers. +// Distributed under the MIT/X11 software license, see the accompanying +// file COPYING or http://www.opensource.org/licenses/mit-license.php. +#ifndef TRIANGLES_CHAINDB_MIGRATE_H +#define TRIANGLES_CHAINDB_MIGRATE_H + +#include + +// Migrate legacy LevelDB chain state from /txleveldb to RocksDB in +// /rocksdb. The source is never modified. Returns true when migration +// succeeds or when there is nothing to migrate. +bool MaybeMigrateLevelDbToRocksDb(bool fForce, std::string& strError); + +#endif // TRIANGLES_CHAINDB_MIGRATE_H diff --git a/src/init.cpp b/src/init.cpp index 586496d..995608b 100644 --- a/src/init.cpp +++ b/src/init.cpp @@ -24,6 +24,7 @@ #endif #include "notificationqueue.h" #include "addressindex.h" +#include "chaindb_migrate.h" #include #include #include @@ -1022,6 +1023,16 @@ bool AppInit2() } } + // ********************************************************* Step 6d: optional LevelDB -> RocksDB chain DB migration + if (GetBoolArg("-migratechaindb", false) || GetBoolArg("-migratechaindbforce", false)) + { + uiInterface.InitMessage(_("Migrating chain database to RocksDB...")); + std::string strMigrateError; + bool fForce = GetBoolArg("-migratechaindbforce", false); + if (!MaybeMigrateLevelDbToRocksDb(fForce, strMigrateError)) + return InitError(strprintf(_("Chain DB migration failed: %s"), strMigrateError.c_str())); + } + // ********************************************************* Step 7: load blockchain if (!bitdb.Open(GetDataDir())) diff --git a/src/main.cpp b/src/main.cpp index e6339ca..85aa516 100644 --- a/src/main.cpp +++ b/src/main.cpp @@ -21,6 +21,7 @@ #include "notificationqueue.h" #include "addressindex.h" #include "snapshotnet.h" +#include "syncmanager.h" #include #include #include @@ -118,36 +119,9 @@ extern enum Checkpoints::CPMode CheckpointsMode; namespace { -struct CHeaderSyncNode -{ - CBlock header; - int nHeight; - uint256 nChainTrust; - bool fRequested; - int64_t nLastRequestTime; - int64_t nFirstRequestTime; // when this block was first requested (for latency tracking) - int64_t nInsertTime; -}; - -static std::map mapHeaderSync; -static uint256 hashBestHeaderSync = 0; -static int64_t nLastNewHeaderTime = 0; static CCriticalSection cs_PostIbdWork; static bool fPostIbdWorkStarted = false; -static const unsigned int MAX_HEADER_SYNC_CACHE = 15000; -static const unsigned int HEADER_DOWNLOAD_WINDOW = 1024; // Wider pipeline for multi-peer parallel IBD -static const unsigned int HEADER_DOWNLOAD_PER_PEER = 32; // Reduced from 64 for Tor circuit stability -static const size_t HEADER_REDUNDANT_PEER_THRESHOLD = 4; // Only do dual-peer redundancy when peer count is below this -static const unsigned int HEADER_SYNC_LOW_WATER = HEADER_DOWNLOAD_WINDOW / 4; -static const unsigned int HEADER_SYNC_TARGET_INFLIGHT = HEADER_DOWNLOAD_WINDOW / 2; -static const int64_t HEADER_REQUEST_TIMEOUT_MICROS = 60 * 1000000; // 60s for Tor latency (was 30s) -static const int64_t HEADER_REDUNDANT_REQUEST_MICROS = 5 * 1000000; // 5s redundant request (reduced for Tor) -static const int64_t HEADER_SYNC_TTL_MICROS = 15 * 60 * 1000000; // 15-minute TTL for cache entries (extended for Tor latency) -static const int64_t HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS = 5; -static const int64_t HEADER_SYNC_CONTROL_INTERVAL_SECONDS = 5; -static const int64_t HEADER_SYNC_WATCHDOG_SECONDS = 25; - static void ThreadPostIbdWork(void* parg) { RenameThread("Triangles-postibd"); @@ -195,499 +169,6 @@ static void ThreadPostIbdWork(void* parg) } } -static uint256 GetHeaderSyncTrust(unsigned int nBits) -{ - CBigNum bnTarget; - bnTarget.SetCompact(nBits); - - if (bnTarget <= 0) - return 0; - - return ((CBigNum(1) << 256) / (bnTarget + 1)).getuint256(); -} - -static bool GetKnownHeaderState(const uint256& hash, int& nHeight, uint256& nChainTrust) -{ - if (auto miBlock = mapBlockIndex.find(hash); miBlock != mapBlockIndex.end()) - { - nHeight = miBlock->second->nHeight; - nChainTrust = miBlock->second->nChainTrust; - return true; - } - - if (auto miHeader = mapHeaderSync.find(hash); miHeader != mapHeaderSync.end()) - { - nHeight = miHeader->second.nHeight; - nChainTrust = miHeader->second.nChainTrust; - return true; - } - - return false; -} - -static bool GetHeaderSyncPrevHash(const uint256& hash, uint256& hashPrev) -{ - if (auto miHeader = mapHeaderSync.find(hash); miHeader != mapHeaderSync.end()) - { - hashPrev = miHeader->second.header.hashPrevBlock; - return true; - } - - if (auto miBlock = mapBlockIndex.find(hash); miBlock != mapBlockIndex.end() && miBlock->second->pprev) - { - hashPrev = miBlock->second->pprev->GetBlockHash(); - return true; - } - - return false; -} - -static void RecomputeBestHeaderSync() -{ - hashBestHeaderSync = 0; - uint256 nBestTrust = 0; - - for (const auto& [hash, node] : mapHeaderSync) - { - if (hashBestHeaderSync == 0 || node.nChainTrust > nBestTrust) - { - hashBestHeaderSync = hash; - nBestTrust = node.nChainTrust; - } - } -} - -static void PruneHeaderSync() -{ - const int64_t nNow = GetTime() * 1000000; - - // TTL eviction: remove entries older than 5 minutes - if (mapHeaderSync.size() > MAX_HEADER_SYNC_CACHE / 2) - { - unsigned int nEvicted = 0; - for (auto it = mapHeaderSync.begin(); it != mapHeaderSync.end(); ) - { - if (nNow - it->second.nInsertTime >= HEADER_SYNC_TTL_MICROS) - { - it = mapHeaderSync.erase(it); - ++nEvicted; - } - else - ++it; - } - if (nEvicted > 0) - { - printf("IBD-DIAG: TTL-evicted %u stale header sync entries, %u remain\n", - nEvicted, (unsigned int)mapHeaderSync.size()); - RecomputeBestHeaderSync(); - } - } - - // Hard limit: if still over max, evict oldest entries - if (mapHeaderSync.size() > MAX_HEADER_SYNC_CACHE) - { - printf("IBD-DIAG: header sync cache exceeded %u entries, evicting oldest\n", MAX_HEADER_SYNC_CACHE); - while (mapHeaderSync.size() > MAX_HEADER_SYNC_CACHE * 3 / 4) - { - // Find oldest entry by insert time - auto oldest = mapHeaderSync.begin(); - for (auto it = mapHeaderSync.begin(); it != mapHeaderSync.end(); ++it) - { - if (it->second.nInsertTime < oldest->second.nInsertTime) - oldest = it; - } - mapHeaderSync.erase(oldest); - } - RecomputeBestHeaderSync(); - } -} - -static bool AddHeaderSyncNode(const CBlock& header, const uint256& hashHeader) -{ - if (mapBlockIndex.count(hashHeader) || mapHeaderSync.count(hashHeader)) - return true; - - if (!header.vtx.empty()) - { - printf("IBD-DIAG: header rejected (has vtx) hash=%s\n", hashHeader.ToString().substr(0,20).c_str()); - return false; - } - - if (header.GetBlockTime() > GetTime() + 15 * 60) - { - printf("IBD-DIAG: header rejected (future time) hash=%s time=%u\n", - hashHeader.ToString().substr(0,20).c_str(), header.nTime); - return false; - } - - int nPrevHeight = -1; - uint256 nPrevChainTrust = 0; - if (!GetKnownHeaderState(header.hashPrevBlock, nPrevHeight, nPrevChainTrust)) - { - printf("IBD-DIAG: header rejected (prev unknown) hash=%s prevHash=%s\n", - hashHeader.ToString().substr(0,20).c_str(), - header.hashPrevBlock.ToString().substr(0,20).c_str()); - return false; - } - - const int nHeight = nPrevHeight + 1; - if (nHeight <= CUTOFF_POW_BLOCK && !CheckProofOfWork(hashHeader, header.nBits)) - { - printf("IBD-DIAG: header PoW FAILED at height %d hash=%s nBits=%08x prevHash=%s\n", - nHeight, hashHeader.ToString().substr(0,20).c_str(), header.nBits, - header.hashPrevBlock.ToString().substr(0,20).c_str()); - return false; - } - - CHeaderSyncNode node; - node.header = header; - node.nHeight = nHeight; - node.nChainTrust = nPrevChainTrust + GetHeaderSyncTrust(header.nBits); - node.fRequested = false; - node.nLastRequestTime = 0; - node.nFirstRequestTime = 0; - node.nInsertTime = GetTime() * 1000000; - - mapHeaderSync.insert({hashHeader, node}); - - if (hashBestHeaderSync == 0 || node.nChainTrust > mapHeaderSync[hashBestHeaderSync].nChainTrust) - hashBestHeaderSync = hashHeader; - - PruneHeaderSync(); - return true; -} - -static CBlockLocator BuildHeaderSyncLocator(uint256 hashTip) -{ - if (hashTip == 0) - return CBlockLocator(pindexBest); - - std::vector vHave; - int nStep = 1; - - while (hashTip != 0) - { - vHave.push_back(hashTip); - - for (int i = 0; i < nStep && hashTip != 0; ++i) - { - uint256 hashPrev = 0; - if (!GetHeaderSyncPrevHash(hashTip, hashPrev)) - hashTip = 0; - else - hashTip = hashPrev; - } - - if (vHave.size() > 10) - nStep *= 2; - } - - vHave.push_back(!fTestNet ? hashGenesisBlockOfficial : hashGenesisBlockTestNet); - return CBlockLocator(vHave); -} - -static std::vector GetHeaderSyncDownloadPath(uint256 hashTip) -{ - std::vector vPath; - - while (hashTip != 0 && !mapBlockIndex.count(hashTip)) - { - auto mi = mapHeaderSync.find(hashTip); - if (mi == mapHeaderSync.end()) - break; - - vPath.push_back(hashTip); - hashTip = mi->second.header.hashPrevBlock; - } - - std::reverse(vPath.begin(), vPath.end()); - return vPath; -} - -static unsigned int CountHeaderSyncInFlight() -{ - const int64_t nNow = GetTime() * 1000000; - unsigned int nInFlight = 0; - for (const auto& [hash, node] : mapHeaderSync) - { - if (node.fRequested && nNow - node.nLastRequestTime < HEADER_REQUEST_TIMEOUT_MICROS) - ++nInFlight; - } - return nInFlight; -} - -static unsigned int GetHeaderSyncPlannerDepth() -{ - if (hashBestHeaderSync == 0) - return 0; - - return (unsigned int)GetHeaderSyncDownloadPath(hashBestHeaderSync).size(); -} - -static int GetHeaderSyncPlannerHeight() -{ - if (hashBestHeaderSync == 0) - return pindexBest ? pindexBest->nHeight : -1; - - auto mi = mapHeaderSync.find(hashBestHeaderSync); - if (mi == mapHeaderSync.end()) - return pindexBest ? pindexBest->nHeight : -1; - - return mi->second.nHeight; -} - -static unsigned int QueueHeaderSyncBlocks(CNode* pfrom, unsigned int nWindow) -{ - if (!pfrom || hashBestHeaderSync == 0) - return 0; - - const std::vector vPath = GetHeaderSyncDownloadPath(hashBestHeaderSync); - if (vPath.empty()) - return 0; - - const int64_t nNow = GetTime() * 1000000; - unsigned int nInFlight = CountHeaderSyncInFlight(); - unsigned int nQueued = 0; - - for (const auto& hash : vPath) - { - if (nInFlight + nQueued >= nWindow) - break; - - auto mi = mapHeaderSync.find(hash); - if (mi == mapHeaderSync.end()) - continue; - - if (mi->second.fRequested && nNow - mi->second.nLastRequestTime < HEADER_REQUEST_TIMEOUT_MICROS) - continue; - - pfrom->AskFor(CInv(MSG_BLOCK, hash)); - mi->second.fRequested = true; - mi->second.nLastRequestTime = nNow; - ++nQueued; - } - - return nQueued; -} - -// Returns the first request time (microseconds) for a block in the header sync cache, or 0 -static int64_t GetHeaderSyncRequestTime(const uint256& hashBlock) -{ - auto mi = mapHeaderSync.find(hashBlock); - if (mi == mapHeaderSync.end()) - return 0; - return mi->second.nFirstRequestTime; -} - -static void MarkHeaderSyncBlockAccepted(const uint256& hashBlock) -{ - auto mi = mapHeaderSync.find(hashBlock); - if (mi == mapHeaderSync.end()) - return; - - mapHeaderSync.erase(mi); - if (hashBestHeaderSync == hashBlock) - RecomputeBestHeaderSync(); -} - -static void ContinueHeaderSync(CNode* pfrom, const uint256& hashTip) -{ - if (!pfrom || hashTip == 0) - return; - - CBlockLocator locator = BuildHeaderSyncLocator(hashTip); - if (locator.IsNull()) - return; - - pfrom->PushMessage("getheaders", locator, uint256(0)); -} - -static bool RequestHeaderSyncRefill(CNode* pfrom, uint256 hashTip, int64_t nMinIntervalSeconds, const char* pszReason) -{ - if (!pfrom || pfrom->fClient || pfrom->nVersion == 0 || !IsInitialBlockDownload()) - return false; - - const int64_t nNowSec = GetTime(); - if (nMinIntervalSeconds > 0 && - nNowSec - pfrom->nLastIbdHeaderRequest < nMinIntervalSeconds) - return false; - - uint256 hashLocatorTip = hashTip; - if (hashLocatorTip == 0 || - (!mapBlockIndex.count(hashLocatorTip) && !mapHeaderSync.count(hashLocatorTip))) - { - hashLocatorTip = hashBestHeaderSync; - } - - if (hashLocatorTip != 0 && (!pindexBest || hashLocatorTip != pindexBest->GetBlockHash())) - { - ContinueHeaderSync(pfrom, hashLocatorTip); - } - else - { - if (!pindexBest) - return false; - - pfrom->pindexLastGetHeadersBegin = nullptr; - pfrom->PushGetHeaders(pindexBest, uint256(0)); - hashLocatorTip = pindexBest->GetBlockHash(); - } - - pfrom->nLastIbdHeaderRequest = nNowSec; - printf("IBD-DIAG: %s getheaders to peer=%s locator=%s plannerDepth=%u inflight=%u\n", - pszReason, pfrom->addr.ToString().c_str(), - hashLocatorTip.ToString().substr(0,20).c_str(), - GetHeaderSyncPlannerDepth(), CountHeaderSyncInFlight()); - return true; -} - -static unsigned int RequestHeaderSyncRefillAllPeers(uint256 hashTip, int64_t nMinIntervalSeconds, const char* pszReason) -{ - std::vector vEligiblePeers; - { - LOCK(cs_vNodes); - for (CNode* pnode : vNodes) - { - if (!pnode->fClient && pnode->nVersion != 0 && !pnode->fDisconnect) - vEligiblePeers.push_back(pnode); - } - } - - unsigned int nRequested = 0; - for (CNode* pnode : vEligiblePeers) - { - if (RequestHeaderSyncRefill(pnode, hashTip, nMinIntervalSeconds, pszReason)) - ++nRequested; - } - - return nRequested; -} - -// Parallel block downloading: distribute blocks across all available peers -static unsigned int QueueHeaderSyncBlocksParallel(unsigned int nWindow) -{ - if (hashBestHeaderSync == 0) - return 0; - - const std::vector vPath = GetHeaderSyncDownloadPath(hashBestHeaderSync); - if (vPath.empty()) - return 0; - - // Collect eligible peers - std::vector vEligiblePeers; - { - LOCK(cs_vNodes); - for (CNode* pnode : vNodes) - { - if (!pnode->fClient && pnode->nVersion != 0 && !pnode->fDisconnect) - vEligiblePeers.push_back(pnode); - } - } - - if (vEligiblePeers.empty()) - return 0; - - const int64_t nNow = GetTime() * 1000000; - unsigned int nInFlight = CountHeaderSyncInFlight(); - unsigned int nQueued = 0; - unsigned int nPeerIndex = 0; - - // Sort peers by blocks delivered (descending) for speed-weighted assignment. - // Faster peers get more blocks assigned to them, improving IBD throughput - // on Tor networks where latency varies significantly between peers. - std::sort(vEligiblePeers.begin(), vEligiblePeers.end(), - [](const CNode* a, const CNode* b) { - return a->nBlocksDelivered > b->nBlocksDelivered; - }); - - // Build a weighted distribution: top peer gets 3 slots per round, second gets 2, rest get 1. - std::vector vWeightedPeers; - for (size_t i = 0; i < vEligiblePeers.size(); i++) - { - int nWeight = (i == 0) ? 3 : (i == 1) ? 2 : 1; - for (int w = 0; w < nWeight; w++) - vWeightedPeers.push_back(vEligiblePeers[i]); - } - - // Adaptive timeout: use average peer latency to set timeouts. - // If peers average 2s, timeout at 10s. If peers average 15s, timeout at 45s. - // Clamp between 10s and 60s. Default to 60s when no latency data. - int64_t nAdaptiveTimeout = HEADER_REQUEST_TIMEOUT_MICROS; - { - int64_t nTotalLatency = 0; - int nPeersWithLatency = 0; - for (const CNode* pnode : vEligiblePeers) { - if (pnode->nAvgBlockLatencyUs > 0) { - nTotalLatency += pnode->nAvgBlockLatencyUs; - ++nPeersWithLatency; - } - } - if (nPeersWithLatency > 0) { - int64_t nAvgLatency = nTotalLatency / nPeersWithLatency; - nAdaptiveTimeout = std::max((int64_t)(10 * 1000000), - std::min((int64_t)(60 * 1000000), nAvgLatency * 5)); - } - } - - // Distribute blocks across peers using speed-weighted assignment - for (const auto& hash : vPath) - { - if (nInFlight + nQueued >= nWindow) - break; - - auto mi = mapHeaderSync.find(hash); - if (mi == mapHeaderSync.end()) - continue; - - bool fNeedsRequest = false; - if (!mi->second.fRequested) - { - fNeedsRequest = true; - } - else if (nNow - mi->second.nLastRequestTime >= nAdaptiveTimeout) - { - fNeedsRequest = true; - } - else if (nNow - mi->second.nLastRequestTime >= HEADER_REDUNDANT_REQUEST_MICROS) - { - fNeedsRequest = true; - } - - if (!fNeedsRequest) - continue; - - CNode* pnode = vWeightedPeers[nPeerIndex % vWeightedPeers.size()]; - pnode->AskFor(CInv(MSG_BLOCK, hash)); - - if (IsInitialBlockDownload() && - vWeightedPeers.size() >= 2 && - vWeightedPeers.size() < HEADER_REDUNDANT_PEER_THRESHOLD && - !mi->second.fRequested) - { - CNode* pnode2 = vWeightedPeers[(nPeerIndex + 1) % vWeightedPeers.size()]; - if (pnode2 != pnode) - pnode2->AskFor(CInv(MSG_BLOCK, hash)); - } - - if (!mi->second.fRequested || nNow - mi->second.nLastRequestTime >= HEADER_REQUEST_TIMEOUT_MICROS) - { - if (!mi->second.fRequested) - mi->second.nFirstRequestTime = nNow; - mi->second.fRequested = true; - mi->second.nLastRequestTime = nNow; - } - - ++nQueued; - ++nPeerIndex; - } - - if (nQueued > 0) - printf("IBD-DIAG: parallel queue distributed %u blocks across %zu peers (window=%u, inflight=%u)\n", - nQueued, vEligiblePeers.size(), nWindow, nInFlight); - - return nQueued; -} - } // namespace ////////////////////////////////////////////////////////////////////////////// @@ -3672,7 +3153,7 @@ bool ProcessBlock(CNode* pfrom, CBlock* pblock) if (!pblock->AcceptBlock()) return error("ProcessBlock() : AcceptBlock FAILED"); - MarkHeaderSyncBlockAccepted(hash); + g_syncManager.BlockAccepted(hash); // Recursively process any orphan blocks that depended on this one vector vWorkQueue; @@ -3688,7 +3169,7 @@ bool ProcessBlock(CNode* pfrom, CBlock* pblock) if (pblockOrphan->AcceptBlock()) { vWorkQueue.push_back(pblockOrphan->GetHash()); - MarkHeaderSyncBlockAccepted(pblockOrphan->GetHash()); + g_syncManager.BlockAccepted(pblockOrphan->GetHash()); } mapOrphanBlocks.erase(pblockOrphan->GetHash()); setStakeSeenOrphan.erase(pblockOrphan->GetProofOfStake()); @@ -3702,16 +3183,16 @@ bool ProcessBlock(CNode* pfrom, CBlock* pblock) if (IsInitialBlockDownload()) { const unsigned int nQueued = - (hashBestHeaderSync != 0) ? QueueHeaderSyncBlocksParallel(HEADER_DOWNLOAD_WINDOW) : 0; + (g_syncManager.GetBestHeader() != 0) ? g_syncManager.QueueBlocksParallel() : 0; if (nQueued > 0) printf("IBD-DIAG: queued %u more blocks from header planner after accepting %s\n", nQueued, hash.ToString().substr(0,20).c_str()); - const unsigned int nPlannerDepth = GetHeaderSyncPlannerDepth(); - if (nPlannerDepth <= HEADER_SYNC_LOW_WATER) + const unsigned int nPlannerDepth = g_syncManager.GetPlannerDepth(); + if (nPlannerDepth <= CSyncManager::HEADER_SYNC_LOW_WATER) { - const unsigned int nRefilled = RequestHeaderSyncRefillAllPeers( - hashBestHeaderSync, HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, + const unsigned int nRefilled = g_syncManager.RequestRefillAllPeers( + g_syncManager.GetBestHeader(), CSyncManager::HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, (nPlannerDepth == 0) ? "post-accept planner empty" : "post-accept planner low-water"); if (nRefilled > 0) printf("IBD-DIAG: post-accept requested headers from %u peers at plannerDepth=%u after block %s\n", @@ -4582,7 +4063,7 @@ bool static ProcessMessage(CNode* pfrom, string strCommand, CDataStream& vRecv) // so we learn about future blocks much faster. The headers handler // will AskFor each unknown block, pre-populating the download queue. if (fIBD) - RequestHeaderSyncRefill(pfrom, hashBestHeaderSync, 0, "version bootstrap"); + g_syncManager.RequestRefill(pfrom, g_syncManager.GetBestHeader(), 0, "version bootstrap"); printf("IBD-DIAG: sent getblocks%s from height %d to peer %s\n", fIBD ? "+getheaders" : "", nBestHeight, pfrom->addr.ToString().c_str()); } @@ -4972,97 +4453,8 @@ bool static ProcessMessage(CNode* pfrom, string strCommand, CDataStream& vRecv) { vector vHeaders; vRecv >> vHeaders; - if (vHeaders.size() > 2000) - { - pfrom->Misbehaving(20); - return error("message headers size() = %" PRIszu "", vHeaders.size()); - } - - uint256 hashChainTip = 0; - int nNewHeaders = 0; - for (const CBlock& header : vHeaders) - { - if (!header.vtx.empty()) - { - pfrom->Misbehaving(20); - return error("headers message includes transactions"); - } - - const uint256 hashHeader = header.GetHash(); - if (mapBlockIndex.count(hashHeader) || mapHeaderSync.count(hashHeader)) - { - hashChainTip = hashHeader; - continue; - } - - if (hashChainTip != 0) - { - if (header.hashPrevBlock != hashChainTip) - { - pfrom->Misbehaving(20); - return error("non-continuous headers sequence"); - } - } - else - { - auto miPrev = mapBlockIndex.find(header.hashPrevBlock); - if (miPrev == mapBlockIndex.end() && !mapHeaderSync.count(header.hashPrevBlock)) - break; - } - - if (!AddHeaderSyncNode(header, hashHeader)) - { - pfrom->Misbehaving(20); - return error("invalid header sequence"); - } - - hashChainTip = hashHeader; - nNewHeaders++; - } - - int nRequested = 0; - if (hashBestHeaderSync != 0) - nRequested = QueueHeaderSyncBlocksParallel(HEADER_DOWNLOAD_WINDOW); - - if (nNewHeaders > 0) - nLastNewHeaderTime = GetTime(); - - if (nNewHeaders > 0 || nRequested > 0) - printf("IBD-DIAG: accepted %d new headers, queued %d blocks from %zu headers (peer=%s bestHeader=%s)\n", - nNewHeaders, nRequested, vHeaders.size(), pfrom->addr.ToString().c_str(), - hashBestHeaderSync.ToString().substr(0,20).c_str()); - - // If we received a full batch, continue fetching headers. - // During IBD, prefer getheaders over getblocks since headers are ~80 bytes - // vs full blocks, letting us discover the chain structure faster. - if (vHeaders.size() >= 2000) - { - if (IsInitialBlockDownload() && hashChainTip != 0) - ContinueHeaderSync(pfrom, hashChainTip); - else - pfrom->PushGetBlocks(pindexBest, uint256(0)); - } - else if (IsInitialBlockDownload() && nNewHeaders > 0 && hashChainTip != 0) - { - // Partial batch with new content. The peer either truncated its - // response (e.g. send-buffer pressure on Tor) or is briefly at the - // tip of what it knows. Either way the v5.9.2 fix only refilled - // when the cache fully drained, so a partial batch could leave the - // pipeline silently parked. Ask this peer to continue from the - // highest header we now know — covers truncated responses, and - // costs at most one empty headers reply when the peer is honestly - // at the chain tip. - ContinueHeaderSync(pfrom, hashChainTip); - } - else if (IsInitialBlockDownload()) - { - const unsigned int nPlannerDepth = GetHeaderSyncPlannerDepth(); - if (nPlannerDepth <= HEADER_SYNC_LOW_WATER) - RequestHeaderSyncRefill( - pfrom, (hashChainTip != 0) ? hashChainTip : hashBestHeaderSync, - HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, - (nPlannerDepth == 0) ? "headers planner empty" : "headers planner low-water"); - } + if (!g_syncManager.ProcessHeaders(pfrom, vHeaders)) + return false; } @@ -5154,24 +4546,7 @@ bool static ProcessMessage(CNode* pfrom, string strCommand, CDataStream& vRecv) CInv inv(MSG_BLOCK, hashBlock); pfrom->AddInventoryKnown(inv); - // Track block delivery and measure latency for adaptive timeouts - pfrom->nBlocksDelivered++; - if (nBestHeight > pfrom->nBestKnownHeight) - pfrom->nBestKnownHeight = nBestHeight; - - // Update rolling average latency (exponential moving average, 7/8 old + 1/8 new) - { - int64_t nRequestTime = GetHeaderSyncRequestTime(hashBlock); - if (nRequestTime > 0) { - int64_t nLatency = GetTime() * 1000000 - nRequestTime; - if (nLatency > 0) { - if (pfrom->nAvgBlockLatencyUs == 0) - pfrom->nAvgBlockLatencyUs = nLatency; - else - pfrom->nAvgBlockLatencyUs = (pfrom->nAvgBlockLatencyUs * 7 + nLatency) / 8; - } - } - } + g_syncManager.TrackBlockDelivery(pfrom, hashBlock); if (ProcessBlock(pfrom, &block)) { @@ -5180,7 +4555,7 @@ bool static ProcessMessage(CNode* pfrom, string strCommand, CDataStream& vRecv) if (IsInitialBlockDownload()) { // Keep download window full after every accepted block - QueueHeaderSyncBlocksParallel(HEADER_DOWNLOAD_WINDOW); + g_syncManager.QueueBlocksParallel(); static int nBlocksSinceRequest = 0; if (++nBlocksSinceRequest >= 500) @@ -5982,18 +5357,18 @@ bool SendMessages(CNode* pto, bool fSendTrickle) { pto->pindexLastGetHeadersBegin = nullptr; - uint256 hashLocatorTip = hashBestHeaderSync; + uint256 hashLocatorTip = g_syncManager.GetBestHeader(); if (hashLocatorTip == 0 && nHighestInvWalk > nBestHeight && hashHighestInvWalk != 0 && mapBlockIndex.count(hashHighestInvWalk)) { hashLocatorTip = hashHighestInvWalk; } - unsigned int nRefilled = RequestHeaderSyncRefillAllPeers( + unsigned int nRefilled = g_syncManager.RequestRefillAllPeers( hashLocatorTip, 0, "stall-recovery"); - unsigned int nQueued = QueueHeaderSyncBlocksParallel(HEADER_DOWNLOAD_WINDOW); + unsigned int nQueued = g_syncManager.QueueBlocksParallel(); printf("SYNC-DIAG: stall recovery used headers-first path (locator=%s, refillPeers=%u, queued=%u)\n", hashLocatorTip.ToString().substr(0,20).c_str(), @@ -6024,84 +5399,9 @@ bool SendMessages(CNode* pto, bool fSendTrickle) } } - // Per-peer IBD getheaders heartbeat. The v5.9.2 belt-and-suspenders - // had three holes that this replaces: - // 1. It only fired when hashBestHeaderSync == 0 (cache fully empty); - // a few stale in-flight entries blocked refill until 15-min TTL. - // 2. The throttle was process-wide, so an unresponsive peer could - // "absorb" the one-per-30s request and leave others unkicked. - // 3. It had no path for "peer stopped feeding mid-batch" — only - // total-cache-drain triggered it. - // - // Per-peer heartbeat with an adaptive interval covers all three: - // - Low-water mode (cache below the download window): 15s, refills - // before the planner runs dry without waiting for cache exhaustion. - // - Steady mode (cache filled): 60s, keeps each peer's view of our - // locator fresh so a peer that goes silent gets re-asked, and a - // peer that catches up between calls can announce new headers. - // An empty headers response is ~14 bytes — cheap on Tor, no abuse risk. - if (!pto->fClient && pto->nVersion != 0 && IsInitialBlockDownload()) - { - const int64_t nNowSec = GetTime(); - const unsigned int nPlannerDepth = GetHeaderSyncPlannerDepth(); - const unsigned int nInFlight = CountHeaderSyncInFlight(); - static int64_t nLastHeaderPlannerControl = 0; - static int64_t nLastHeaderWatchdog = 0; - static int64_t nLastBlockPlannerControl = 0; - - if (nLastNewHeaderTime == 0) - nLastNewHeaderTime = nNowSec; - - if (nNowSec - nLastHeaderPlannerControl >= HEADER_SYNC_CONTROL_INTERVAL_SECONDS && - nPlannerDepth < HEADER_SYNC_LOW_WATER && - nInFlight < HEADER_SYNC_TARGET_INFLIGHT) - { - const unsigned int nRefilled = RequestHeaderSyncRefillAllPeers( - hashBestHeaderSync, HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, - "control-loop"); - if (nRefilled > 0) - printf("IBD-DIAG: control-loop refill from %u peers (plannerDepth=%u inflight=%u target=%u)\n", - nRefilled, nPlannerDepth, nInFlight, HEADER_SYNC_TARGET_INFLIGHT); - nLastHeaderPlannerControl = nNowSec; - } - - if (nNowSec - nLastHeaderWatchdog >= HEADER_SYNC_CONTROL_INTERVAL_SECONDS && - nNowSec - nLastNewHeaderTime >= HEADER_SYNC_WATCHDOG_SECONDS) - { - const unsigned int nRefilled = RequestHeaderSyncRefillAllPeers( - hashBestHeaderSync, HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, - "headers-watchdog"); - if (nRefilled > 0) - printf("IBD-DIAG: headers watchdog refill from %u peers after %llds without new headers (plannerDepth=%u inflight=%u)\n", - nRefilled, - (long long)(nNowSec - nLastNewHeaderTime), - nPlannerDepth, - nInFlight); - nLastHeaderWatchdog = nNowSec; - } - - const int64_t nMinInterval = - (mapHeaderSync.size() < HEADER_DOWNLOAD_WINDOW) ? 15 : 60; - - if (nNowSec - pto->nLastIbdHeaderRequest >= nMinInterval) - RequestHeaderSyncRefill(pto, hashBestHeaderSync, nMinInterval, "heartbeat"); - - // Keep the block planner alive even when no new headers arrive and - // no blocks are being accepted. Without this periodic kick, the - // redundant-request and timeout logic inside QueueHeaderSyncBlocksParallel() - // only runs on header arrivals or block acceptance, so IBD can park - // indefinitely behind one missing frontier block. - if (nNowSec - nLastBlockPlannerControl >= HEADER_SYNC_CONTROL_INTERVAL_SECONDS && - hashBestHeaderSync != 0 && - nPlannerDepth > 0) - { - const unsigned int nRequeued = QueueHeaderSyncBlocksParallel(HEADER_DOWNLOAD_WINDOW); - if (nRequeued > 0) - printf("IBD-DIAG: block-planner control queued %u block requests (plannerDepth=%u inflight=%u)\n", - nRequeued, nPlannerDepth, nInFlight); - nLastBlockPlannerControl = nNowSec; - } - } + // Per-peer IBD getheaders heartbeat and block-planner cadence. + // Logic lives in CSyncManager::Tick — see syncmanager.cpp. + g_syncManager.Tick(pto, nHighestInvWalk, hashHighestInvWalk); // // Message: getdata @@ -6112,9 +5412,9 @@ bool SendMessages(CNode* pto, bool fSendTrickle) if (GetTime() - nLastStatus >= 15) { printf("IBD-DIAG: STATUS height=%d plannerHeight=%d plannerDepth=%u inflight=%u peers=%d askfor_queued=%d orphans=%d\n", nBestHeight, - GetHeaderSyncPlannerHeight(), - GetHeaderSyncPlannerDepth(), - CountHeaderSyncInFlight(), + g_syncManager.GetPlannerHeight(), + g_syncManager.GetPlannerDepth(), + g_syncManager.CountInFlight(), (int)vNodes.size(), (int)pto->mapAskFor.size(), (int)mapOrphanBlocks.size()); diff --git a/src/syncmanager.cpp b/src/syncmanager.cpp new file mode 100644 index 0000000..f06837b --- /dev/null +++ b/src/syncmanager.cpp @@ -0,0 +1,664 @@ +// Copyright (c) 2026 The Triangles developers +// Distributed under the MIT/X11 software license, see the accompanying +// file COPYING or http://www.opensource.org/licenses/mit-license.php. + +#include "syncmanager.h" + +#include "bignum.h" +#include "checkpoints.h" +#include "main.h" +#include "net.h" +#include "util.h" + +#include +#include + +struct CSyncManager::HeaderNode +{ + CBlock header; + int nHeight; + uint256 nChainTrust; + bool fRequested; + int64_t nLastRequestTime; + int64_t nFirstRequestTime; + int64_t nInsertTime; +}; + +namespace +{ +static const unsigned int MAX_HEADER_SYNC_CACHE = 15000; +static const size_t HEADER_REDUNDANT_PEER_THRESHOLD = 4; +static const int64_t HEADER_REQUEST_TIMEOUT_MICROS = 60 * 1000000; +static const int64_t HEADER_REDUNDANT_REQUEST_MICROS = 5 * 1000000; +static const int64_t HEADER_SYNC_TTL_MICROS = 15 * 60 * 1000000; + +std::map mapHeaders; +uint256 hashBestHeader = 0; +int64_t nLastNewHeaderTime = 0; +} + +CSyncManager g_syncManager; + +bool CSyncManager::HaveHeader(const uint256& hash) const +{ + return mapHeaders.count(hash) != 0; +} + +uint256 CSyncManager::GetBestHeader() const +{ + return hashBestHeader; +} + +std::size_t CSyncManager::GetHeaderCount() const +{ + return mapHeaders.size(); +} + +uint256 CSyncManager::GetHeaderTrust(unsigned int nBits) const +{ + CBigNum bnTarget; + bnTarget.SetCompact(nBits); + + if (bnTarget <= 0) + return 0; + + return ((CBigNum(1) << 256) / (bnTarget + 1)).getuint256(); +} + +bool CSyncManager::GetKnownHeaderState(const uint256& hash, int& nHeight, uint256& nChainTrust) const +{ + std::map::const_iterator miBlock = mapBlockIndex.find(hash); + if (miBlock != mapBlockIndex.end()) + { + nHeight = miBlock->second->nHeight; + nChainTrust = miBlock->second->nChainTrust; + return true; + } + + std::map::const_iterator miHeader = mapHeaders.find(hash); + if (miHeader != mapHeaders.end()) + { + nHeight = miHeader->second.nHeight; + nChainTrust = miHeader->second.nChainTrust; + return true; + } + + return false; +} + +bool CSyncManager::GetPrevHash(const uint256& hash, uint256& hashPrev) const +{ + std::map::const_iterator miHeader = mapHeaders.find(hash); + if (miHeader != mapHeaders.end()) + { + hashPrev = miHeader->second.header.hashPrevBlock; + return true; + } + + std::map::const_iterator miBlock = mapBlockIndex.find(hash); + if (miBlock != mapBlockIndex.end() && miBlock->second->pprev) + { + hashPrev = miBlock->second->pprev->GetBlockHash(); + return true; + } + + return false; +} + +void CSyncManager::RecomputeBestHeader() +{ + hashBestHeader = 0; + uint256 nBestTrust = 0; + + for (std::map::const_iterator it = mapHeaders.begin(); it != mapHeaders.end(); ++it) + { + if (hashBestHeader == 0 || it->second.nChainTrust > nBestTrust) + { + hashBestHeader = it->first; + nBestTrust = it->second.nChainTrust; + } + } +} + +void CSyncManager::PruneHeaders() +{ + const int64_t nNow = GetTime() * 1000000; + + if (mapHeaders.size() > MAX_HEADER_SYNC_CACHE / 2) + { + unsigned int nEvicted = 0; + for (std::map::iterator it = mapHeaders.begin(); it != mapHeaders.end(); ) + { + if (nNow - it->second.nInsertTime >= HEADER_SYNC_TTL_MICROS) + { + it = mapHeaders.erase(it); + ++nEvicted; + } + else + ++it; + } + if (nEvicted > 0) + { + printf("IBD-DIAG: TTL-evicted %u stale sync headers, %u remain\n", + nEvicted, (unsigned int)mapHeaders.size()); + RecomputeBestHeader(); + } + } + + if (mapHeaders.size() > MAX_HEADER_SYNC_CACHE) + { + printf("IBD-DIAG: sync header cache exceeded %u entries, evicting oldest\n", MAX_HEADER_SYNC_CACHE); + while (mapHeaders.size() > MAX_HEADER_SYNC_CACHE * 3 / 4) + { + std::map::iterator oldest = mapHeaders.begin(); + for (std::map::iterator it = mapHeaders.begin(); it != mapHeaders.end(); ++it) + { + if (it->second.nInsertTime < oldest->second.nInsertTime) + oldest = it; + } + mapHeaders.erase(oldest); + } + RecomputeBestHeader(); + } +} + +bool CSyncManager::AddHeaderNode(const CBlock& header, const uint256& hashHeader) +{ + if (mapBlockIndex.count(hashHeader) || mapHeaders.count(hashHeader)) + return true; + + if (!header.vtx.empty()) + { + printf("IBD-DIAG: header rejected (has vtx) hash=%s\n", hashHeader.ToString().substr(0,20).c_str()); + return false; + } + + if (header.GetBlockTime() > GetTime() + 15 * 60) + { + printf("IBD-DIAG: header rejected (future time) hash=%s time=%u\n", + hashHeader.ToString().substr(0,20).c_str(), header.nTime); + return false; + } + + int nPrevHeight = -1; + uint256 nPrevChainTrust = 0; + if (!GetKnownHeaderState(header.hashPrevBlock, nPrevHeight, nPrevChainTrust)) + { + printf("IBD-DIAG: header rejected (prev unknown) hash=%s prevHash=%s\n", + hashHeader.ToString().substr(0,20).c_str(), + header.hashPrevBlock.ToString().substr(0,20).c_str()); + return false; + } + + const int nHeight = nPrevHeight + 1; + if (nHeight <= CUTOFF_POW_BLOCK && !CheckProofOfWork(hashHeader, header.nBits)) + { + printf("IBD-DIAG: header PoW FAILED at height %d hash=%s nBits=%08x prevHash=%s\n", + nHeight, hashHeader.ToString().substr(0,20).c_str(), header.nBits, + header.hashPrevBlock.ToString().substr(0,20).c_str()); + return false; + } + + HeaderNode node; + node.header = header; + node.nHeight = nHeight; + node.nChainTrust = nPrevChainTrust + GetHeaderTrust(header.nBits); + node.fRequested = false; + node.nLastRequestTime = 0; + node.nFirstRequestTime = 0; + node.nInsertTime = GetTime() * 1000000; + + mapHeaders.insert({hashHeader, node}); + + if (hashBestHeader == 0 || node.nChainTrust > mapHeaders[hashBestHeader].nChainTrust) + hashBestHeader = hashHeader; + + PruneHeaders(); + return true; +} + +std::vector CSyncManager::GetDownloadPath(uint256 hashTip) const +{ + std::vector vPath; + + while (hashTip != 0 && !mapBlockIndex.count(hashTip)) + { + std::map::const_iterator mi = mapHeaders.find(hashTip); + if (mi == mapHeaders.end()) + break; + + vPath.push_back(hashTip); + hashTip = mi->second.header.hashPrevBlock; + } + + std::reverse(vPath.begin(), vPath.end()); + return vPath; +} + +unsigned int CSyncManager::CountInFlight() const +{ + const int64_t nNow = GetTime() * 1000000; + unsigned int nInFlight = 0; + for (std::map::const_iterator it = mapHeaders.begin(); it != mapHeaders.end(); ++it) + { + if (it->second.fRequested && nNow - it->second.nLastRequestTime < HEADER_REQUEST_TIMEOUT_MICROS) + ++nInFlight; + } + return nInFlight; +} + +unsigned int CSyncManager::GetPlannerDepth() const +{ + if (hashBestHeader == 0) + return 0; + + return (unsigned int)GetDownloadPath(hashBestHeader).size(); +} + +int CSyncManager::GetPlannerHeight() const +{ + if (hashBestHeader == 0) + return pindexBest ? pindexBest->nHeight : -1; + + std::map::const_iterator mi = mapHeaders.find(hashBestHeader); + if (mi == mapHeaders.end()) + return pindexBest ? pindexBest->nHeight : -1; + + return mi->second.nHeight; +} + +int64_t CSyncManager::GetRequestTime(const uint256& hashBlock) const +{ + std::map::const_iterator mi = mapHeaders.find(hashBlock); + if (mi == mapHeaders.end()) + return 0; + return mi->second.nFirstRequestTime; +} + +void CSyncManager::BlockAccepted(const uint256& hashBlock) +{ + std::map::iterator mi = mapHeaders.find(hashBlock); + if (mi == mapHeaders.end()) + return; + + mapHeaders.erase(mi); + if (hashBestHeader == hashBlock) + RecomputeBestHeader(); +} + +void CSyncManager::ContinueHeaders(CNode* pfrom, const uint256& hashTip) +{ + if (!pfrom || hashTip == 0) + return; + + std::vector vHave; + uint256 hashWalk = hashTip; + int nStep = 1; + + while (hashWalk != 0) + { + vHave.push_back(hashWalk); + + for (int i = 0; i < nStep && hashWalk != 0; ++i) + { + uint256 hashPrev = 0; + if (!GetPrevHash(hashWalk, hashPrev)) + hashWalk = 0; + else + hashWalk = hashPrev; + } + + if (vHave.size() > 10) + nStep *= 2; + } + + vHave.push_back(!fTestNet ? hashGenesisBlockOfficial : hashGenesisBlockTestNet); + pfrom->PushMessage("getheaders", CBlockLocator(vHave), uint256(0)); +} + +bool CSyncManager::RequestRefill(CNode* pfrom, uint256 hashTip, int64_t nMinIntervalSeconds, const char* pszReason) +{ + if (!pfrom || pfrom->fClient || pfrom->nVersion == 0 || !IsInitialBlockDownload()) + return false; + + const int64_t nNowSec = GetTime(); + if (nMinIntervalSeconds > 0 && + nNowSec - pfrom->nLastIbdHeaderRequest < nMinIntervalSeconds) + return false; + + uint256 hashLocatorTip = hashTip; + if (hashLocatorTip == 0 || + (!mapBlockIndex.count(hashLocatorTip) && !mapHeaders.count(hashLocatorTip))) + { + hashLocatorTip = hashBestHeader; + } + + if (hashLocatorTip != 0 && (!pindexBest || hashLocatorTip != pindexBest->GetBlockHash())) + { + ContinueHeaders(pfrom, hashLocatorTip); + } + else + { + if (!pindexBest) + return false; + + pfrom->pindexLastGetHeadersBegin = NULL; + pfrom->PushGetHeaders(pindexBest, uint256(0)); + hashLocatorTip = pindexBest->GetBlockHash(); + } + + pfrom->nLastIbdHeaderRequest = nNowSec; + printf("IBD-DIAG: %s getheaders to peer=%s locator=%s plannerDepth=%u inflight=%u\n", + pszReason, pfrom->addr.ToString().c_str(), + hashLocatorTip.ToString().substr(0,20).c_str(), + GetPlannerDepth(), CountInFlight()); + return true; +} + +unsigned int CSyncManager::RequestRefillAllPeers(uint256 hashTip, int64_t nMinIntervalSeconds, const char* pszReason) +{ + std::vector vEligiblePeers; + { + LOCK(cs_vNodes); + for (CNode* pnode : vNodes) + { + if (!pnode->fClient && pnode->nVersion != 0 && !pnode->fDisconnect) + vEligiblePeers.push_back(pnode); + } + } + + unsigned int nRequested = 0; + for (CNode* pnode : vEligiblePeers) + { + if (RequestRefill(pnode, hashTip, nMinIntervalSeconds, pszReason)) + ++nRequested; + } + + return nRequested; +} + +unsigned int CSyncManager::QueueBlocksParallel(unsigned int nWindow) +{ + if (hashBestHeader == 0) + return 0; + + const std::vector vPath = GetDownloadPath(hashBestHeader); + if (vPath.empty()) + return 0; + + std::vector vEligiblePeers; + { + LOCK(cs_vNodes); + for (CNode* pnode : vNodes) + { + if (!pnode->fClient && pnode->nVersion != 0 && !pnode->fDisconnect) + vEligiblePeers.push_back(pnode); + } + } + + if (vEligiblePeers.empty()) + return 0; + + const int64_t nNow = GetTime() * 1000000; + unsigned int nInFlight = CountInFlight(); + unsigned int nQueued = 0; + unsigned int nPeerIndex = 0; + + std::sort(vEligiblePeers.begin(), vEligiblePeers.end(), + [](const CNode* a, const CNode* b) { + return a->nBlocksDelivered > b->nBlocksDelivered; + }); + + std::vector vWeightedPeers; + for (size_t i = 0; i < vEligiblePeers.size(); i++) + { + int nWeight = (i == 0) ? 3 : (i == 1) ? 2 : 1; + for (int w = 0; w < nWeight; w++) + vWeightedPeers.push_back(vEligiblePeers[i]); + } + + int64_t nAdaptiveTimeout = HEADER_REQUEST_TIMEOUT_MICROS; + { + int64_t nTotalLatency = 0; + int nPeersWithLatency = 0; + for (const CNode* pnode : vEligiblePeers) + { + if (pnode->nAvgBlockLatencyUs > 0) + { + nTotalLatency += pnode->nAvgBlockLatencyUs; + ++nPeersWithLatency; + } + } + if (nPeersWithLatency > 0) + { + int64_t nAvgLatency = nTotalLatency / nPeersWithLatency; + nAdaptiveTimeout = std::max((int64_t)(10 * 1000000), + std::min((int64_t)(60 * 1000000), nAvgLatency * 5)); + } + } + + for (std::vector::const_iterator it = vPath.begin(); it != vPath.end(); ++it) + { + if (nInFlight + nQueued >= nWindow) + break; + + std::map::iterator mi = mapHeaders.find(*it); + if (mi == mapHeaders.end()) + continue; + + bool fNeedsRequest = false; + if (!mi->second.fRequested) + fNeedsRequest = true; + else if (nNow - mi->second.nLastRequestTime >= nAdaptiveTimeout) + fNeedsRequest = true; + else if (nNow - mi->second.nLastRequestTime >= HEADER_REDUNDANT_REQUEST_MICROS) + fNeedsRequest = true; + + if (!fNeedsRequest) + continue; + + CNode* pnode = vWeightedPeers[nPeerIndex % vWeightedPeers.size()]; + pnode->AskFor(CInv(MSG_BLOCK, *it)); + + if (IsInitialBlockDownload() && + vWeightedPeers.size() >= 2 && + vWeightedPeers.size() < HEADER_REDUNDANT_PEER_THRESHOLD && + !mi->second.fRequested) + { + CNode* pnode2 = vWeightedPeers[(nPeerIndex + 1) % vWeightedPeers.size()]; + if (pnode2 != pnode) + pnode2->AskFor(CInv(MSG_BLOCK, *it)); + } + + if (!mi->second.fRequested || nNow - mi->second.nLastRequestTime >= HEADER_REQUEST_TIMEOUT_MICROS) + { + if (!mi->second.fRequested) + mi->second.nFirstRequestTime = nNow; + mi->second.fRequested = true; + mi->second.nLastRequestTime = nNow; + } + + ++nQueued; + ++nPeerIndex; + } + + if (nQueued > 0) + printf("IBD-DIAG: sync manager queued %u blocks across %zu peers (window=%u, inflight=%u)\n", + nQueued, vEligiblePeers.size(), nWindow, nInFlight); + + return nQueued; +} + +bool CSyncManager::ProcessHeaders(CNode* pfrom, const std::vector& vHeaders) +{ + if (vHeaders.size() > 2000) + { + pfrom->Misbehaving(20); + return error("message headers size() = %" PRIszu "", vHeaders.size()); + } + + uint256 hashChainTip = 0; + int nNewHeaders = 0; + for (const CBlock& header : vHeaders) + { + if (!header.vtx.empty()) + { + pfrom->Misbehaving(20); + return error("headers message includes transactions"); + } + + const uint256 hashHeader = header.GetHash(); + if (mapBlockIndex.count(hashHeader) || mapHeaders.count(hashHeader)) + { + hashChainTip = hashHeader; + continue; + } + + if (hashChainTip != 0) + { + if (header.hashPrevBlock != hashChainTip) + { + pfrom->Misbehaving(20); + return error("non-continuous headers sequence"); + } + } + else + { + std::map::iterator miPrev = mapBlockIndex.find(header.hashPrevBlock); + if (miPrev == mapBlockIndex.end() && !mapHeaders.count(header.hashPrevBlock)) + break; + } + + if (!AddHeaderNode(header, hashHeader)) + { + pfrom->Misbehaving(20); + return error("invalid header sequence"); + } + + hashChainTip = hashHeader; + nNewHeaders++; + } + + int nRequested = 0; + if (hashBestHeader != 0) + nRequested = QueueBlocksParallel(HEADER_DOWNLOAD_WINDOW); + + if (nNewHeaders > 0) + nLastNewHeaderTime = GetTime(); + + if (nNewHeaders > 0 || nRequested > 0) + printf("IBD-DIAG: accepted %d new headers, queued %d blocks from %zu headers (peer=%s bestHeader=%s)\n", + nNewHeaders, nRequested, vHeaders.size(), pfrom->addr.ToString().c_str(), + hashBestHeader.ToString().substr(0,20).c_str()); + + if (vHeaders.size() >= 2000) + { + if (IsInitialBlockDownload() && hashChainTip != 0) + ContinueHeaders(pfrom, hashChainTip); + else + pfrom->PushGetBlocks(pindexBest, uint256(0)); + } + else if (IsInitialBlockDownload() && nNewHeaders > 0 && hashChainTip != 0) + { + ContinueHeaders(pfrom, hashChainTip); + } + else if (IsInitialBlockDownload()) + { + const unsigned int nPlannerDepth = GetPlannerDepth(); + if (nPlannerDepth <= HEADER_SYNC_LOW_WATER) + RequestRefill( + pfrom, (hashChainTip != 0) ? hashChainTip : hashBestHeader, + HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, + (nPlannerDepth == 0) ? "headers planner empty" : "headers planner low-water"); + } + + return true; +} + +void CSyncManager::TrackBlockDelivery(CNode* pfrom, const uint256& hashBlock) +{ + if (!pfrom) + return; + + pfrom->nBlocksDelivered++; + if (nBestHeight > pfrom->nBestKnownHeight) + pfrom->nBestKnownHeight = nBestHeight; + + int64_t nRequestTime = GetRequestTime(hashBlock); + if (nRequestTime > 0) + { + int64_t nLatency = GetTime() * 1000000 - nRequestTime; + if (nLatency > 0) + { + if (pfrom->nAvgBlockLatencyUs == 0) + pfrom->nAvgBlockLatencyUs = nLatency; + else + pfrom->nAvgBlockLatencyUs = (pfrom->nAvgBlockLatencyUs * 7 + nLatency) / 8; + } + } +} + +void CSyncManager::Tick(CNode* pto, int nHighestInvWalk, const uint256& hashHighestInvWalk) +{ + if (!pto || pto->fClient || pto->nVersion == 0 || !IsInitialBlockDownload()) + return; + + const int64_t nNowSec = GetTime(); + const unsigned int nPlannerDepth = GetPlannerDepth(); + const unsigned int nInFlight = CountInFlight(); + static int64_t nLastHeaderPlannerControl = 0; + static int64_t nLastHeaderWatchdog = 0; + static int64_t nLastBlockPlannerControl = 0; + + if (nLastNewHeaderTime == 0) + nLastNewHeaderTime = nNowSec; + + if (nNowSec - nLastHeaderPlannerControl >= HEADER_SYNC_CONTROL_INTERVAL_SECONDS && + nPlannerDepth < HEADER_SYNC_LOW_WATER && + nInFlight < HEADER_SYNC_TARGET_INFLIGHT) + { + const unsigned int nRefilled = RequestRefillAllPeers( + hashBestHeader, HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, + "control-loop"); + if (nRefilled > 0) + printf("IBD-DIAG: control-loop refill from %u peers (plannerDepth=%u inflight=%u target=%u)\n", + nRefilled, nPlannerDepth, nInFlight, HEADER_SYNC_TARGET_INFLIGHT); + nLastHeaderPlannerControl = nNowSec; + } + + if (nNowSec - nLastHeaderWatchdog >= HEADER_SYNC_CONTROL_INTERVAL_SECONDS && + nNowSec - nLastNewHeaderTime >= HEADER_SYNC_WATCHDOG_SECONDS) + { + const unsigned int nRefilled = RequestRefillAllPeers( + hashBestHeader, HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, + "headers-watchdog"); + if (nRefilled > 0) + printf("IBD-DIAG: headers watchdog refill from %u peers after %llds without new headers (plannerDepth=%u inflight=%u)\n", + nRefilled, + (long long)(nNowSec - nLastNewHeaderTime), + nPlannerDepth, + nInFlight); + nLastHeaderWatchdog = nNowSec; + } + + const int64_t nMinInterval = (mapHeaders.size() < HEADER_DOWNLOAD_WINDOW) ? 15 : 60; + if (nNowSec - pto->nLastIbdHeaderRequest >= nMinInterval) + RequestRefill(pto, hashBestHeader, nMinInterval, "heartbeat"); + + if (nNowSec - nLastBlockPlannerControl >= HEADER_SYNC_CONTROL_INTERVAL_SECONDS && + hashBestHeader != 0 && + nPlannerDepth > 0) + { + const unsigned int nRequeued = QueueBlocksParallel(HEADER_DOWNLOAD_WINDOW); + if (nRequeued > 0) + printf("IBD-DIAG: block-planner control queued %u block requests (plannerDepth=%u inflight=%u)\n", + nRequeued, nPlannerDepth, nInFlight); + nLastBlockPlannerControl = nNowSec; + } + + if (hashBestHeader == 0 && nHighestInvWalk > nBestHeight && + hashHighestInvWalk != 0 && mapBlockIndex.count(hashHighestInvWalk)) + { + RequestRefill(pto, hashHighestInvWalk, HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS, "inv-walk bridge"); + } +} diff --git a/src/syncmanager.h b/src/syncmanager.h new file mode 100644 index 0000000..03cb837 --- /dev/null +++ b/src/syncmanager.h @@ -0,0 +1,58 @@ +// Copyright (c) 2026 The Triangles developers +// Distributed under the MIT/X11 software license, see the accompanying +// file COPYING or http://www.opensource.org/licenses/mit-license.php. +#ifndef TRIANGLES_SYNCMANAGER_H +#define TRIANGLES_SYNCMANAGER_H + +#include "uint256.h" + +#include +#include +#include + +class CBlock; +class CInv; +class CNode; + +class CSyncManager +{ +public: + struct HeaderNode; + + static constexpr unsigned int HEADER_DOWNLOAD_WINDOW = 1024; + static constexpr unsigned int HEADER_SYNC_LOW_WATER = HEADER_DOWNLOAD_WINDOW / 4; + static constexpr unsigned int HEADER_SYNC_TARGET_INFLIGHT = HEADER_DOWNLOAD_WINDOW / 2; + static constexpr int64_t HEADER_SYNC_REFILL_MIN_INTERVAL_SECONDS = 5; + static constexpr int64_t HEADER_SYNC_CONTROL_INTERVAL_SECONDS = 5; + static constexpr int64_t HEADER_SYNC_WATCHDOG_SECONDS = 25; + + bool HaveHeader(const uint256& hash) const; + uint256 GetBestHeader() const; + std::size_t GetHeaderCount() const; + unsigned int CountInFlight() const; + unsigned int GetPlannerDepth() const; + int GetPlannerHeight() const; + int64_t GetRequestTime(const uint256& hashBlock) const; + + bool RequestRefill(CNode* pfrom, uint256 hashTip, int64_t nMinIntervalSeconds, const char* pszReason); + unsigned int RequestRefillAllPeers(uint256 hashTip, int64_t nMinIntervalSeconds, const char* pszReason); + unsigned int QueueBlocksParallel(unsigned int nWindow = HEADER_DOWNLOAD_WINDOW); + bool ProcessHeaders(CNode* pfrom, const std::vector& vHeaders); + void BlockAccepted(const uint256& hashBlock); + void TrackBlockDelivery(CNode* pfrom, const uint256& hashBlock); + void Tick(CNode* pto, int nHighestInvWalk, const uint256& hashHighestInvWalk); + +private: + uint256 GetHeaderTrust(unsigned int nBits) const; + bool GetKnownHeaderState(const uint256& hash, int& nHeight, uint256& nChainTrust) const; + bool GetPrevHash(const uint256& hash, uint256& hashPrev) const; + void RecomputeBestHeader(); + void PruneHeaders(); + bool AddHeaderNode(const CBlock& header, const uint256& hashHeader); + std::vector GetDownloadPath(uint256 hashTip) const; + void ContinueHeaders(CNode* pfrom, const uint256& hashTip); +}; + +extern CSyncManager g_syncManager; + +#endif // TRIANGLES_SYNCMANAGER_H diff --git a/src/txdb-rocksdb.h b/src/txdb-rocksdb.h index e4a168f..ac5b66d 100644 --- a/src/txdb-rocksdb.h +++ b/src/txdb-rocksdb.h @@ -38,6 +38,15 @@ public: bool LoadBlockIndex() override; + // Write a raw serialized key/value pair, bypassing the typed Write<>() + // overloads. Intended for the chaindb migration utility, which carries + // bytes directly across from a CTxDB (LevelDB) iterator. Honors the + // active write batch if one is open. + bool WriteRawRecordForMigration(const std::string& key, const std::string& value) + { + return WriteRaw(key, value); + } + protected: bool ReadRaw(const std::string& key, std::string& value) const override; bool WriteRaw(const std::string& key, const std::string& value) override;