diff --git a/blockchain/chain.go b/blockchain/chain.go index d1d984b9..f047c664 100644 --- a/blockchain/chain.go +++ b/blockchain/chain.go @@ -348,6 +348,19 @@ func (b *BlockChain) DisableVerify(disable bool) { b.chainLock.Unlock() } +// HaveHeader returns whether or not the chain instance has the block header +// represented by the passed hash. Note that this will return true for both the +// main chain and any side chains. +// +// This function is safe for concurrent access. +func (b *BlockChain) HaveHeader(hash *chainhash.Hash) bool { + b.index.RLock() + node := b.index.index[*hash] + headerKnown := node != nil + b.index.RUnlock() + return headerKnown +} + // HaveBlock returns whether or not the chain instance has the block represented // by the passed hash. This includes checking the various places a block can // be like part of the main chain or on a side chain. @@ -1436,7 +1449,7 @@ func (b *BlockChain) isOldTimestamp(node *blockNode) bool { } // maybeUpdateIsCurrent potentially updates whether or not the chain believes it -// is current. +// is current using the provided best chain tip. // // It makes use of a latching approach such that once the chain becomes current // it will only switch back to false in the case no new blocks have been seen @@ -1474,7 +1487,20 @@ func (b *BlockChain) maybeUpdateIsCurrent(curBest *blockNode) { log.Debugf("Chain latched to current at block %s (height %d)", curBest.hash, curBest.height) } +} +// MaybeUpdateIsCurrent potentially updates whether or not the chain believes it +// is current. +// +// It makes use of a latching approach such that once the chain becomes current +// it will only switch back to false in the case no new blocks have been seen +// for an extended period of time. +// +// This function is safe for concurrent access. +func (b *BlockChain) MaybeUpdateIsCurrent() { + b.chainLock.Lock() + b.maybeUpdateIsCurrent(b.bestChain.Tip()) + b.chainLock.Unlock() } // isCurrent returns whether or not the chain believes it is current based on @@ -2189,9 +2215,8 @@ func New(ctx context.Context, config *Config) (*BlockChain, error) { b.dbInfo.bidxVer) tip := b.bestChain.Tip() - log.Infof("Chain state: height %d, hash %v, total transactions %d, "+ - "work %v, stake version %v", tip.height, tip.hash, - b.stateSnapshot.TotalTxns, tip.workSum, 0) + log.Infof("Chain state: height %d, hash %v, total transactions %d, work %v", + tip.height, tip.hash, b.stateSnapshot.TotalTxns, tip.workSum) return &b, nil } diff --git a/blockchain/chainquery.go b/blockchain/chainquery.go index dd3ffe33..1ce5d03e 100644 --- a/blockchain/chainquery.go +++ b/blockchain/chainquery.go @@ -1,4 +1,4 @@ -// Copyright (c) 2018-2020 The Decred developers +// Copyright (c) 2018-2021 The Decred developers // Use of this source code is governed by an ISC // license that can be found in the LICENSE file. @@ -154,3 +154,82 @@ func (b *BlockChain) BestInvalidHeader() chainhash.Hash { b.index.RUnlock() return hash } + +// NextNeededBlocks returns hashes for the next blocks after the current best +// chain tip that are needed to make progress towards the current best known +// header skipping any blocks that are already known or in the provided map of +// blocks to exclude. Typically the caller would want to exclude all blocks +// that have outstanding requests. +// +// The maximum number of results is limited to the provided value or the maximum +// allowed by the internal lookahead buffer in the case the requested number of +// max results exceeds that value. +// +// This function is safe for concurrent access. +func (b *BlockChain) NextNeededBlocks(maxResults uint8, exclude map[chainhash.Hash]struct{}) []*chainhash.Hash { + // Nothing to do when no results are requested. + if maxResults == 0 { + return nil + } + + // Determine the common ancestor between the current best chain tip and the + // current best known header. In practice this should never be nil because + // orphan headers are not allowed into the block index, but be paranoid and + // check anyway in case things change in the future. + b.index.RLock() + bestHeader := b.index.bestHeader + fork := b.bestChain.FindFork(bestHeader) + if fork == nil { + b.index.RUnlock() + return nil + } + + // Determine the final block to consider for determining the next needed + // blocks by determining the descendants of the current best chain tip on + // the branch that leads to the best known header while clamping the number + // of descendants to consider to a lookahead buffer. + const lookaheadBuffer = 512 + numBlocksToConsider := bestHeader.height - fork.height + if numBlocksToConsider == 0 { + b.index.RUnlock() + return nil + } + if numBlocksToConsider > lookaheadBuffer { + bestHeader = bestHeader.Ancestor(fork.height + lookaheadBuffer) + numBlocksToConsider = lookaheadBuffer + } + + // Walk backwards from the final block to consider to the current best chain + // tip excluding any blocks that already have their data available or that + // the caller asked to be excluded (likely because they've already been + // requested). + neededBlocks := make([]*chainhash.Hash, 0, numBlocksToConsider) + for node := bestHeader; node != nil && node != fork; node = node.parent { + _, isExcluded := exclude[node.hash] + if isExcluded || node.status.HaveData() { + continue + } + + neededBlocks = append(neededBlocks, &node.hash) + } + b.index.RUnlock() + + // Reverse the needed blocks so they are in forwards order. + reverse := func(s []*chainhash.Hash) { + slen := len(s) + for i := 0; i < slen/2; i++ { + s[i], s[slen-1-i] = s[slen-1-i], s[i] + } + } + reverse(neededBlocks) + + // Clamp the number of results to the lower of the requested max or number + // available. + if int64(maxResults) > numBlocksToConsider { + maxResults = uint8(numBlocksToConsider) + } + if uint16(len(neededBlocks)) > uint16(maxResults) { + neededBlocks = neededBlocks[:maxResults] + } + return neededBlocks +} diff --git a/internal/netsync/manager.go b/internal/netsync/manager.go index 982fbd06..cda1fbf2 100644 --- a/internal/netsync/manager.go +++ b/internal/netsync/manager.go @@ -6,10 +6,10 @@ package netsync import ( - "container/list" "context" "errors" "fmt" + "math" "runtime/debug" "sync" "time" @@ -30,10 +30,13 @@ import ( const ( // minInFlightBlocks is the minimum number of blocks that should be - // in the request queue for headers-first mode before requesting - // more. + // in the request queue before requesting more. minInFlightBlocks = 10 + // maxInFlightBlocks is the maximum number of blocks to allow in the sync + // peer request queue. + maxInFlightBlocks = 16 + // maxRejectedTxns is the maximum number of rejected transactions // hashes to store in memory. maxRejectedTxns = 1000 @@ -45,6 +48,21 @@ const ( // maxRequestedTxns is the maximum number of requested transactions // hashes to store in memory. maxRequestedTxns = wire.MaxInvPerMsg + + // maxExpectedHeaderAnnouncementsPerMsg is the maximum number of headers in + // a single message that is expected when determining when the message + // appears to be announcing new blocks. + maxExpectedHeaderAnnouncementsPerMsg = 12 + + // maxConsecutiveOrphanHeaders is the maximum number of consecutive header + // messages that contain headers which do not connect a peer can send before + // it is deemed to have diverged so far it is no longer useful. + maxConsecutiveOrphanHeaders = 10 + + // headerSyncStallTimeoutSecs is the number of seconds to wait for progress + // during the header sync process before stalling the sync and disconnecting + // the peer. + headerSyncStallTimeoutSecs = (3 + wire.MaxBlockHeadersPerMsg/1000) * 2 ) // zeroHash is the zero value hash (all zeros). It is defined as a convenience. @@ -158,13 +176,6 @@ type processTransactionMsg struct { reply chan processTransactionResponse } -// headerNode is used as a node in a list of headers that are linked together -// between checkpoints. -type headerNode struct { - height int64 - hash *chainhash.Hash -} - // syncMgrPeer extends a peer to maintain additional state maintained by the // sync manager. type syncMgrPeer struct { @@ -173,6 +184,71 @@ type syncMgrPeer struct { syncCandidate bool requestedTxns map[chainhash.Hash]struct{} requestedBlocks map[chainhash.Hash]struct{} + + // numConsecutiveOrphanHeaders tracks the number of consecutive header + // messages sent by the peer that contain headers which do not connect. It + // is used to detect peers that have either diverged so far they are no + // longer useful or are otherwise being malicious. + numConsecutiveOrphanHeaders int32 +} + +// headerSyncState houses the state used to track the header sync progress and +// related stall handling. +type headerSyncState struct { + // headersSynced tracks whether or not the headers are synced to a point + // that is recent enough to start downloading blocks. + headersSynced bool + + // These fields are used to implement a progress stall timeout that can be + // reset at any time without needing to create a new one and the associated + // extra garbage. + // + // stallTimer is an underlying timer that is used to implement the timeout. + // + // stallChanDrained indicates whether or not the channel for the stall timer + // has already been read and is used when resetting the timer to ensure the + // channel is drained when the timer is stopped as described in the timer + // documentation. + stallTimer *time.Timer + stallChanDrained bool +} + +// makeHeaderSyncState returns a header sync state that is ready to use. +func makeHeaderSyncState() headerSyncState { + stallTimer := time.NewTimer(math.MaxInt64) + stallTimer.Stop() + return headerSyncState{ + stallTimer: stallTimer, + stallChanDrained: true, + } +} + +// stopStallTimeout stops the progress stall timer while ensuring to read from +// the timer's channel in the case the timer already expired which can happen +// due to the fact the stop happens in between channel reads. This behavior is +// well documented in the Timer docs. +// +// NOTE: This function must not be called concurrent with any other receives on +// the timer's channel. +func (state *headerSyncState) stopStallTimeout() { + t := state.stallTimer + if !t.Stop() && !state.stallChanDrained { + <-t.C + } + state.stallChanDrained = true +} + +// resetStallTimeout resets the progress stall timer while ensuring to read from +// the timer's channel in the case the timer already expired which can happen +// due to the fact the reset happens in between channel reads. This behavior is +// well documented in the Timer docs. +// +// NOTE: This function must not be called concurrent with any other receives on +// the timer's channel. +func (state *headerSyncState) resetStallTimeout() { + state.stopStallTimeout() + state.stallTimer.Reset(headerSyncStallTimeoutSecs * time.Second) + state.stallChanDrained = false } // SyncManager provides a concurrency safe sync manager for handling all @@ -195,11 +271,9 @@ type SyncManager struct { msgChan chan interface{} peers map[*peerpkg.Peer]*syncMgrPeer - // The following fields are used for headers-first mode. - headersFirstMode bool - headerList *list.List - startHeader *list.Element - nextCheckpoint *chaincfg.Checkpoint + // hdrSyncState houses the state used to track the initial header sync + // process and related stall handling. + hdrSyncState headerSyncState // The following fields are used to track the height being synced to from // peers. @@ -225,55 +299,12 @@ func lookupPeer(peer *peerpkg.Peer, peers map[*peerpkg.Peer]*syncMgrPeer) *syncM return sp } -// resetHeaderState sets the headers-first mode state to values appropriate for -// syncing from a new peer. -func (m *SyncManager) resetHeaderState(newestHash *chainhash.Hash, newestHeight int64) { - m.headersFirstMode = false - m.headerList.Init() - m.startHeader = nil - - // When there is a next checkpoint, add an entry for the latest known - // block into the header pool. This allows the next downloaded header - // to prove it links to the chain properly. - if m.nextCheckpoint != nil { - node := headerNode{height: newestHeight, hash: newestHash} - m.headerList.PushBack(&node) - } -} - // SyncHeight returns latest known block being synced to. func (m *SyncManager) SyncHeight() int64 { m.syncHeightMtx.Lock() - defer m.syncHeightMtx.Unlock() - return m.syncHeight -} - -// findNextHeaderCheckpoint returns the next checkpoint after the passed height. -// It returns nil when there is not one either because the height is already -// later than the final checkpoint or some other reason such as disabled -// checkpoints. -func (m *SyncManager) findNextHeaderCheckpoint(height int64) *chaincfg.Checkpoint { - checkpoints := m.cfg.Chain.Checkpoints() - if len(checkpoints) == 0 { - return nil - } - - // There is no next checkpoint if the height is already after the final - // checkpoint. - finalCheckpoint := &checkpoints[len(checkpoints)-1] - if height >= finalCheckpoint.Height { - return nil - } - - // Find the next checkpoint. - nextCheckpoint := finalCheckpoint - for i := len(checkpoints) - 2; i >= 0; i-- { - if height >= checkpoints[i].Height { - break - } - nextCheckpoint = &checkpoints[i] - } - return nextCheckpoint + syncHeight := m.syncHeight + m.syncHeightMtx.Unlock() + return syncHeight } // chainBlockLocatorToHashes converts a block locator from chain to a slice @@ -290,17 +321,61 @@ func chainBlockLocatorToHashes(locator blockchain.BlockLocator) []chainhash.Hash return result } +// fetchNextBlocks creates and sends a request to the provided peer for the next +// blocks to be downloaded based on the current headers. +func (m *SyncManager) fetchNextBlocks(peer *syncMgrPeer) { + // Nothing to do if the target maximum number of blocks to request from the + // peer at the same time are already in flight. + numInFlight := len(peer.requestedBlocks) + if numInFlight >= maxInFlightBlocks { + return + } + + // Determine the next blocks to download based on the final block that has + // already been requested and the next blocks in the branch leading up to + // best known header. + chain := m.cfg.Chain + maxNeeded := uint8(maxInFlightBlocks - numInFlight) + neededBlocks := chain.NextNeededBlocks(maxNeeded, m.requestedBlocks) + if len(neededBlocks) == 0 { + return + } + + // Ensure the number of needed blocks does not exceed the max inventory + // per message. This should never happen because the constants above limit + // it to relatively small values. However, the code below relies on this + // assumption, so assert it. + if len(neededBlocks) > wire.MaxInvPerMsg { + log.Warnf("%d needed blocks exceeds max allowed per message", + len(neededBlocks)) + return + } + + // Build and send a getdata request for needed blocks. + gdmsg := wire.NewMsgGetDataSizeHint(uint(len(neededBlocks))) + for i := range neededBlocks { + hash := neededBlocks[i] + iv := wire.NewInvVect(wire.InvTypeBlock, hash) + m.requestedBlocks[*hash] = struct{}{} + peer.requestedBlocks[*hash] = struct{}{} + gdmsg.AddInvVect(iv) + } + peer.QueueMessage(gdmsg, nil) + +} + // startSync will choose the best peer among the available candidate peers to // download/sync the blockchain from. When syncing is already running, it // simply returns. It also examines the candidates for any which are no longer // candidates and removes them as needed. func (m *SyncManager) startSync() { - // Return now if we're already syncing. + // Nothing more to do when already syncing. if m.syncPeer != nil { return } - best := m.cfg.Chain.BestSnapshot() + chain := m.cfg.Chain + best := chain.BestSnapshot() var bestPeer *syncMgrPeer for _, peer := range m.peers { if !peer.syncCandidate { @@ -318,7 +393,7 @@ func (m *SyncManager) startSync() { continue } - // the best sync candidate is the most updated peer + // The best sync candidate is the most updated peer. if bestPeer == nil { bestPeer = peer } @@ -330,77 +405,78 @@ func (m *SyncManager) startSync() { // Update the state of whether or not the manager believes the chain is // fully synced to whatever the chain believes when there is no candidate // for a sync peer. + // + // Also, return now when there isn't a sync peer candidate as there is + // nothing more to do without one. if bestPeer == nil { m.isCurrentMtx.Lock() - m.isCurrent = m.cfg.Chain.IsCurrent() + m.isCurrent = chain.IsCurrent() + m.isCurrentMtx.Unlock() + log.Warnf("No sync peer candidates available") + return + } + + // Start syncing from the best peer. + + // Clear the requestedBlocks if the sync peer changes, otherwise + // we may ignore blocks we need that the last sync peer failed + // to send. + m.requestedBlocks = make(map[chainhash.Hash]struct{}) + + syncHeight := bestPeer.LastBlock() + + headersSynced := m.hdrSyncState.headersSynced + if !headersSynced { + log.Infof("Syncing headers to block height %d from peer %v", syncHeight, + bestPeer) + } + + // The chain is not synced whenever the current best height is less than the + // height to sync to. + if best.Height < syncHeight { + m.isCurrentMtx.Lock() + m.isCurrent = false m.isCurrentMtx.Unlock() } - // Start syncing from the best peer if one was selected. - if bestPeer != nil { - // Clear the requestedBlocks if the sync peer changes, otherwise - // we may ignore blocks we need that the last sync peer failed - // to send. - m.requestedBlocks = make(map[chainhash.Hash]struct{}) + // Request headers to discover any blocks that are not already known + // starting from the parent of the best known header for the local chain. + // The parent is used as a means to accurately discover the best known block + // of the remote peer in the case both tips are the same where it would + // otherwise result in an empty response. + bestHeaderHash, _ := chain.BestHeader() + parentHash := bestHeaderHash + header, err := chain.HeaderByHash(&bestHeaderHash) + if err == nil { + parentHash = header.PrevBlock + } + blkLocator := chain.BlockLocatorFromHash(&parentHash) + locator := chainBlockLocatorToHashes(blkLocator) + bestPeer.PushGetHeadersMsg(locator, &zeroHash) - blkLocator := m.cfg.Chain.LatestBlockLocator() - locator := chainBlockLocatorToHashes(blkLocator) + // Track the sync peer and update the sync height when it is higher than the + // currently best known value. + m.syncPeer = bestPeer + m.syncHeightMtx.Lock() + if syncHeight > m.syncHeight { + m.syncHeight = syncHeight + } + m.syncHeightMtx.Unlock() - log.Infof("Syncing to block height %d from peer %v", - bestPeer.LastBlock(), bestPeer.Addr()) + // Start the header sync progress stall timeout when the initial headers + // sync is not already done. + if !headersSynced { + m.hdrSyncState.resetStallTimeout() + } - // The chain is not synced whenever the current best height is less than - // the height to sync to. - if best.Height < bestPeer.LastBlock() { - m.isCurrentMtx.Lock() - m.isCurrent = false - m.isCurrentMtx.Unlock() - } - - // When the current height is less than a known checkpoint we - // can use block headers to learn about which blocks comprise - // the chain up to the checkpoint and perform less validation - // for them. This is possible since each header contains the - // hash of the previous header and a merkle root. Therefore if - // we validate all of the received headers link together - // properly and the checkpoint hashes match, we can be sure the - // hashes for the blocks in between are accurate. Further, once - // the full blocks are downloaded, the merkle root is computed - // and compared against the value in the header which proves the - // full block hasn't been tampered with. - // - // Once we have passed the final checkpoint, or checkpoints are - // disabled, use standard inv messages learn about the blocks - // and fully validate them. Finally, regression test mode does - // not support the headers-first approach so do normal block - // downloads when in regression test mode. - if m.nextCheckpoint != nil && - best.Height < m.nextCheckpoint.Height && - !m.cfg.DisableCheckpoints { - - err := bestPeer.PushGetHeadersMsg(locator, m.nextCheckpoint.Hash) - if err != nil { - log.Errorf("Failed to push getheadermsg for the latest "+ - "blocks: %v", err) - return - } - m.headersFirstMode = true - log.Infof("Downloading headers for blocks %d to %d from peer %s", - best.Height+1, m.nextCheckpoint.Height, bestPeer.Addr()) - } else { - err := bestPeer.PushGetBlocksMsg(locator, &zeroHash) - if err != nil { - log.Errorf("Failed to push getblocksmsg for the latest "+ - "blocks: %v", err) - return - } - } - m.syncPeer = bestPeer - m.syncHeightMtx.Lock() - m.syncHeight = bestPeer.LastBlock() - m.syncHeightMtx.Unlock() - } else { - log.Warnf("No sync peer candidates available") + // Download any blocks needed to catch the local chain up to the best + // known header (if any) when the initial headers sync is already done. + // + // This is done in addition to the header request above to avoid waiting + // for the round trip when there are still blocks that are needed + // regardless of the headers response. + if headersSynced { + m.fetchNextBlocks(m.syncPeer) } } @@ -507,28 +583,23 @@ func (m *SyncManager) handleDonePeerMsg(p *peerpkg.Peer) { delete(m.peers, p) // Remove requested transactions from the global map so that they will - // be fetched from elsewhere next time we get an inv. + // be fetched from elsewhere. for txHash := range peer.requestedTxns { delete(m.requestedTxns, txHash) } // Remove requested blocks from the global map so that they will be - // fetched from elsewhere next time we get an inv. + // fetched from elsewhere. // TODO(oga) we could possibly here check which peers have these blocks // and request them now to speed things up a little. for blockHash := range peer.requestedBlocks { delete(m.requestedBlocks, blockHash) } - // Attempt to find a new peer to sync from if the quitting peer is the - // sync peer. Also, reset the headers-first state if in headers-first - // mode so + // Attempt to find a new peer to sync from and reset the final requested + // block when the quitting peer is the sync peer. if m.syncPeer == peer { m.syncPeer = nil - if m.headersFirstMode { - best := m.cfg.Chain.BestSnapshot() - m.resetHeaderState(&best.Hash, best.Height) - } m.startSync() } } @@ -715,53 +786,33 @@ func (m *SyncManager) handleBlockMsg(bmsg *blockMsg) { blockHash := bmsg.block.Hash() if _, exists := peer.requestedBlocks[*blockHash]; !exists { log.Warnf("Got unrequested block %v from %s -- disconnecting", - blockHash, peer.Addr()) + blockHash, peer) peer.Disconnect() return } - // When in headers-first mode, if the block matches the hash of the - // first header in the list of headers that are being fetched, it's - // eligible for less validation since the headers have already been - // verified to link together and are valid up to the next checkpoint. - // Also, remove the list entry for all blocks except the checkpoint - // since it is needed to verify the next round of headers links - // properly. - isCheckpointBlock := false - behaviorFlags := blockchain.BFNone - if m.headersFirstMode { - firstNodeEl := m.headerList.Front() - if firstNodeEl != nil { - firstNode := firstNodeEl.Value.(*headerNode) - if blockHash.IsEqual(firstNode.hash) { - behaviorFlags |= blockchain.BFFastAdd - if firstNode.hash.IsEqual(m.nextCheckpoint.Hash) { - isCheckpointBlock = true - } else { - m.headerList.Remove(firstNodeEl) - } - } - } - } - - // Remove block from request maps. Either chain will know about it and - // so we shouldn't have any more instances of trying to fetch it, or we - // will fail the insert and thus we'll retry next time we get an inv. + // Process the block to include validation, best chain selection, etc. + // + // Also, remove the block from the request maps once it has been processed. + // This ensures chain is aware of the block before it is removed from the + // maps in order to help prevent duplicate requests. + forkLen, err := m.processBlock(bmsg.block, blockchain.BFNone) delete(peer.requestedBlocks, *blockHash) delete(m.requestedBlocks, *blockHash) - - // Process the block to include validation, best chain selection, etc. - forkLen, err := m.processBlock(bmsg.block, behaviorFlags) if err != nil { - // Ignore orphan blocks. - if errors.Is(err, blockchain.ErrMissingParent) { + // Ideally there should never be any requests for duplicate blocks, but + // ignore any that manage to make it through. + if errors.Is(err, blockchain.ErrDuplicateBlock) { return } // When the error is a rule error, it means the block was simply - // rejected as opposed to something actually going wrong, so log - // it as such. Otherwise, something really did go wrong, so log - // it as an actual error. + // rejected as opposed to something actually going wrong, so log it as + // such. Otherwise, something really did go wrong, so log it as an + // actual error. + // + // Note that orphan blocks are never requested so there is no need to + // test for that rule error separately. var rErr blockchain.RuleError if errors.As(err, &rErr) { log.Infof("Rejected block %v from %s: %v", blockHash, peer, err) @@ -772,8 +823,7 @@ func (m *SyncManager) handleBlockMsg(bmsg *blockMsg) { log.Errorf("Critical failure: %v", err) } - // Convert the error into an appropriate reject message and - // send it. + // Convert the error into an appropriate reject message and send it. code, reason := errToWireRejectCode(err) peer.PushRejectMsg(wire.CmdBlock, code, reason, blockHash, false) return @@ -819,96 +869,10 @@ func (m *SyncManager) handleBlockMsg(bmsg *blockMsg) { } } - // Nothing more to do if we aren't in headers-first mode. - if !m.headersFirstMode { - return - } - - // This is headers-first mode, so if the block is not a checkpoint - // request more blocks using the header list when the request queue is - // getting short. - if !isCheckpointBlock { - if m.startHeader != nil && - len(peer.requestedBlocks) < minInFlightBlocks { - m.fetchHeaderBlocks() - } - return - } - - // This is headers-first mode and the block is a checkpoint. When - // there is a next checkpoint, get the next round of headers by asking - // for headers starting from the block after this one up to the next - // checkpoint. - prevHeight := m.nextCheckpoint.Height - prevHash := m.nextCheckpoint.Hash - m.nextCheckpoint = m.findNextHeaderCheckpoint(prevHeight) - if m.nextCheckpoint != nil { - locator := []chainhash.Hash{*prevHash} - err := peer.PushGetHeadersMsg(locator, m.nextCheckpoint.Hash) - if err != nil { - log.Warnf("Failed to send getheaders message to peer %s: %v", - peer.Addr(), err) - return - } - log.Infof("Downloading headers for blocks %d to %d from peer %s", - prevHeight+1, m.nextCheckpoint.Height, m.syncPeer.Addr()) - return - } - - // This is headers-first mode, the block is a checkpoint, and there are - // no more checkpoints, so switch to normal mode by requesting blocks - // from the block after this one up to the end of the chain (zero hash). - m.headersFirstMode = false - m.headerList.Init() - log.Infof("Reached the final checkpoint -- switching to normal mode") - locator := []chainhash.Hash{*blockHash} - err = peer.PushGetBlocksMsg(locator, &zeroHash) - if err != nil { - log.Warnf("Failed to send getblocks message to peer %s: %v", - peer.Addr(), err) - return - } -} - -// fetchHeaderBlocks creates and sends a request to the syncPeer for the next -// list of blocks to be downloaded based on the current list of headers. -func (m *SyncManager) fetchHeaderBlocks() { - // Nothing to do if there is no start header. - if m.startHeader == nil { - log.Warnf("fetchHeaderBlocks called with no start header") - return - } - - // Build up a getdata request for the list of blocks the headers - // describe. The size hint will be limited to wire.MaxInvPerMsg by - // the function, so no need to double check it here. - gdmsg := wire.NewMsgGetDataSizeHint(uint(m.headerList.Len())) - numRequested := 0 - for e := m.startHeader; e != nil; e = e.Next() { - node, ok := e.Value.(*headerNode) - if !ok { - log.Warn("Header list node type is not a headerNode") - continue - } - - if m.needBlock(node.hash) { - iv := wire.NewInvVect(wire.InvTypeBlock, node.hash) - m.requestedBlocks[*node.hash] = struct{}{} - m.syncPeer.requestedBlocks[*node.hash] = struct{}{} - err := gdmsg.AddInvVect(iv) - if err != nil { - log.Warnf("Failed to add invvect while fetching block "+ - "headers: %v", err) - } - numRequested++ - } - m.startHeader = e.Next() - if numRequested >= wire.MaxInvPerMsg { - break - } - } - if len(gdmsg.InvList) > 0 { - m.syncPeer.QueueMessage(gdmsg, nil) + // Request more blocks using the headers when the request queue is getting + // short. + if peer == m.syncPeer && len(peer.requestedBlocks) < minInFlightBlocks { + m.fetchNextBlocks(peer) } } @@ -919,97 +883,216 @@ func (m *SyncManager) handleHeadersMsg(hmsg *headersMsg) { return } - // The remote peer is misbehaving if we didn't request headers. - msg := hmsg.headers - numHeaders := len(msg.Headers) - if !m.headersFirstMode { - log.Warnf("Got %d unrequested headers from %s -- disconnecting", - numHeaders, peer.Addr()) - peer.Disconnect() - return - } - - // Nothing to do for an empty headers message. + // Nothing to do for an empty headers message as it means the sending peer + // does not have any additional headers for the requested block locator. + headers := hmsg.headers.Headers + numHeaders := len(headers) if numHeaders == 0 { return } - // Process all of the received headers ensuring each one connects to the - // previous and that checkpoints match. - receivedCheckpoint := false - var finalHash *chainhash.Hash - for _, blockHeader := range msg.Headers { - blockHash := blockHeader.BlockHash() - finalHash = &blockHash - - // Ensure there is a previous header to compare against. - prevNodeEl := m.headerList.Back() - if prevNodeEl == nil { - log.Warnf("Header list does not contain a previous element as " + - "expected -- disconnecting peer") - peer.Disconnect() + // Handle the case where the first header does not connect to any known + // headers specially. + chain := m.cfg.Chain + firstHeader := headers[0] + firstHeaderHash := firstHeader.BlockHash() + firstHeaderConnects := chain.HaveHeader(&firstHeader.PrevBlock) + headersSynced := m.hdrSyncState.headersSynced + if !firstHeaderConnects { + // Ignore headers that do not connect to any known headers when the + // initial headers sync is taking place. It is expected that headers + // will be announced that are not yet known. + if !headersSynced { return } - // Ensure the header properly connects to the previous one and - // add it to the list of headers. - node := headerNode{hash: &blockHash} - prevNode := prevNodeEl.Value.(*headerNode) - if prevNode.hash.IsEqual(&blockHeader.PrevBlock) { - node.height = prevNode.height + 1 - e := m.headerList.PushBack(&node) - if m.startHeader == nil { - m.startHeader = e + // Attempt to detect block announcements which do not connect to any + // known headers and request any headers starting from the best header + // the local chain knows in order to (hopefully) discover the missing + // headers. + // + // Meanwhile, also keep track of how many times the peer has + // consecutively sent a headers message that does not connect and + // disconnect it once the max allowed threshold has been reached. + if numHeaders < maxExpectedHeaderAnnouncementsPerMsg { + peer.numConsecutiveOrphanHeaders++ + if peer.numConsecutiveOrphanHeaders >= maxConsecutiveOrphanHeaders { + log.Debugf("Received %d consecutive headers messages that do "+ + "not connect from peer %s -- disconnecting", + peer.numConsecutiveOrphanHeaders, peer) + peer.Disconnect() } - } else { - log.Warnf("Received block header that does not properly connect "+ - "to the chain from peer %s -- disconnecting", peer.Addr()) - peer.Disconnect() + + log.Debugf("Requesting missing parents for header %s (height %d) "+ + "received from peer %s", firstHeaderHash, firstHeader.Height, + peer) + bestHeaderHash, _ := chain.BestHeader() + blkLocator := chain.BlockLocatorFromHash(&bestHeaderHash) + locator := chainBlockLocatorToHashes(blkLocator) + peer.PushGetHeadersMsg(locator, &zeroHash) + return } - // Verify the header at the next checkpoint height matches. - if node.height == m.nextCheckpoint.Height { - if node.hash.IsEqual(m.nextCheckpoint.Hash) { - receivedCheckpoint = true - log.Infof("Verified downloaded block header against "+ - "checkpoint at height %d/hash %s", node.height, node.hash) - } else { - log.Warnf("Block header at height %d/hash %s from peer %s "+ - "does NOT match expected checkpoint hash of %s -- "+ - "disconnecting", node.height, node.hash, peer.Addr(), - m.nextCheckpoint.Hash) + // The initial headers sync process is done and this does not appear to + // be a block announcement, so disconnect the peer. + log.Debugf("Received orphan header from peer %s -- disconnecting", peer) + peer.Disconnect() + return + } + + // Ensure all of the received headers connect the previous one before + // attempting to perform any further processing on any of them. + headerHashes := make([]chainhash.Hash, 0, len(headers)) + headerHashes = append(headerHashes, firstHeaderHash) + for prevIdx, header := range headers[1:] { + prevHash := &headerHashes[prevIdx] + prevHeight := headers[prevIdx].Height + if header.PrevBlock != *prevHash || header.Height != prevHeight+1 { + log.Debugf("Received block header that does not properly connect "+ + "to previous one from peer %s -- disconnecting", peer) + peer.Disconnect() + return + } + headerHashes = append(headerHashes, header.BlockHash()) + } + + // Save the current best known header height prior to processing the headers + // so the code later is able to determine if any new useful headers were + // provided. + _, prevBestHeaderHeight := chain.BestHeader() + + // Process all of the received headers. + for _, header := range headers { + err := chain.ProcessBlockHeader(header, blockchain.BFNone) + if err != nil { + // Note that there is no need to check for an orphan header here + // because they were already verified to connect above. + + log.Debugf("Failed to process block header %s from peer %s: %v -- "+ + "disconnecting", header.BlockHash(), peer, err) + peer.Disconnect() + return + } + } + + // All of the headers were either accepted or already known valid at this + // point. + + // Reset the header sync progress stall timeout when the headers are not + // already synced and progress was made. + newBestHeaderHash, newBestHeaderHeight := chain.BestHeader() + if peer == m.syncPeer && !headersSynced { + if newBestHeaderHeight > prevBestHeaderHeight { + m.hdrSyncState.resetStallTimeout() + } + } + + // Reset the count of consecutive headers messages that contained headers + // which do not connect. Note that this is intentionally only done when all + // of the provided headers are successfully processed above. + peer.numConsecutiveOrphanHeaders = 0 + + // Update the last announced block to the final one in the announced headers + // above and update the height for the peer too. + finalHeader := headers[len(headers)-1] + finalReceivedHash := &headerHashes[len(headerHashes)-1] + peer.UpdateLastAnnouncedBlock(finalReceivedHash) + peer.UpdateLastBlockHeight(int64(finalHeader.Height)) + + // Update the sync height if the new best known header height exceeds it. + syncHeight := m.SyncHeight() + if newBestHeaderHeight > syncHeight { + syncHeight = newBestHeaderHeight + m.syncHeightMtx.Lock() + m.syncHeight = syncHeight + m.syncHeightMtx.Unlock() + } + + // Disconnect outbound peers that have less cumulative work than the minimum + // value already known to have been achieved on the network a priori while + // the initial sync is still underway. This is determined by noting that a + // peer only sends fewer than the maximum number of headers per message when + // it has reached its best known header. + isChainCurrent := chain.IsCurrent() + receivedMaxHeaders := len(headers) == wire.MaxBlockHeadersPerMsg + if !isChainCurrent && !peer.Inbound() && !receivedMaxHeaders { + minKnownWork := m.cfg.ChainParams.MinKnownChainWork + if minKnownWork != nil { + workSum, err := chain.ChainWork(finalReceivedHash) + if err == nil && workSum.Cmp(minKnownWork) < 0 { + log.Debugf("Best known chain for peer %s has too little "+ + "cumulative work -- disconnecting", peer) peer.Disconnect() return } - break } } - // When this header is a checkpoint, switch to fetching the blocks for - // all of the headers since the last checkpoint. - if receivedCheckpoint { - // Since the first entry of the list is always the final block - // that is already in the database and is only used to ensure - // the next header links properly, it must be removed before - // fetching the blocks. - m.headerList.Remove(m.headerList.Front()) - log.Infof("Received %v block headers: Fetching blocks", - m.headerList.Len()) - m.progressLogger.SetLastLogTime(time.Now()) - m.fetchHeaderBlocks() - return + // Request more headers when the peer announced the maximum number of + // headers that can be sent in a single message since it probably has more. + if receivedMaxHeaders { + blkLocator := chain.BlockLocatorFromHash(finalReceivedHash) + locator := chainBlockLocatorToHashes(blkLocator) + peer.PushGetHeadersMsg(locator, &zeroHash) } - // This header is not a checkpoint, so request the next batch of - // headers starting from the latest known header and ending with the - // next checkpoint. - locator := []chainhash.Hash{*finalHash} - err := peer.PushGetHeadersMsg(locator, m.nextCheckpoint.Hash) - if err != nil { - log.Warnf("Failed to send getheaders message to peer %s: %v", - peer.Addr(), err) - return + // Consider the headers synced once the sync peer sends a message with a + // final header that is within a few blocks of the sync height. + if !headersSynced && peer == m.syncPeer { + const syncHeightFetchOffset = 6 + if int64(finalHeader.Height)+syncHeightFetchOffset > syncHeight { + headersSynced = true + m.hdrSyncState.headersSynced = headersSynced + m.hdrSyncState.stopStallTimeout() + + log.Infof("Initial headers sync complete (best header hash %s, "+ + "height %d)", newBestHeaderHash, newBestHeaderHeight) + log.Info("Syncing chain") + m.progressLogger.SetLastLogTime(time.Now()) + + // Potentially update whether the chain believes it is current now + // that the headers are synced. + chain.MaybeUpdateIsCurrent() + isChainCurrent = chain.IsCurrent() + } + } + + // Immediately download blocks associated with the announced headers once + // the chain is current. This allows downloading from whichever peer + // announces it first and also ensures any side chain blocks are downloaded + // for vote consideration. + // + // Ultimately, it would likely be better for this functionality to be moved + // to the code which determines the next blocks to request based on the + // available headers once that code supports downloading from multiple peers + // and associated infrastructure to efficiently determinine which peers have + // the associated block(s). + if isChainCurrent { + gdmsg := wire.NewMsgGetDataSizeHint(uint(len(headers))) + for i := range headerHashes { + // Skip the block when it has already been requested or is otherwise + // already known. + hash := &headerHashes[i] + _, isRequestedBlock := m.requestedBlocks[*hash] + if isRequestedBlock || chain.HaveBlock(hash) { + continue + } + + iv := wire.NewInvVect(wire.InvTypeBlock, hash) + limitAdd(m.requestedBlocks, *hash, maxRequestedBlocks) + limitAdd(peer.requestedBlocks, *hash, maxRequestedBlocks) + gdmsg.AddInvVect(iv) + } + if len(gdmsg.InvList) > 0 { + peer.QueueMessage(gdmsg, nil) + } + } + + // Download any blocks needed to catch the local chain up to the best known + // header (if any) once the initial headers sync is done. + if headersSynced && m.syncPeer != nil { + m.fetchNextBlocks(m.syncPeer) } } @@ -1038,12 +1121,6 @@ func (m *SyncManager) handleNotFoundMsg(nfmsg *notFoundMsg) { } } -// needBlock returns whether or not the block needs to be downloaded. For -// example, it does not need to be downloaded when it is already known. -func (m *SyncManager) needBlock(hash *chainhash.Hash) bool { - return !m.cfg.Chain.HaveBlock(hash) -} - // needTx returns whether or not the transaction needs to be downloaded. For // example, it does not need to be downloaded when it is already known. func (m *SyncManager) needTx(hash *chainhash.Hash) bool { @@ -1066,14 +1143,14 @@ func (m *SyncManager) needTx(hash *chainhash.Hash) bool { return false } -// handleInvMsg handles inv messages from all peers. -// We examine the inventory advertised by the remote peer and act accordingly. +// handleInvMsg handles inv messages from all peers. This entails examining the +// inventory advertised by the remote peer for block and transaction +// announcements and acting accordingly. func (m *SyncManager) handleInvMsg(imsg *invMsg) { peer := lookupPeer(imsg.peer, m.peers) if peer == nil { return } - fromSyncPeer := peer == m.syncPeer isCurrent := m.IsCurrent() // Update state information regarding per-peer known inventory and determine @@ -1087,6 +1164,14 @@ func (m *SyncManager) handleInvMsg(imsg *invMsg) { for _, iv := range imsg.inv.InvList { switch iv.Type { case wire.InvTypeBlock: + // NOTE: All block announcements are now made via headers and the + // decisions regarding which blocks to download are based on those + // headers. Therefore, there is no need to request anything here. + // + // Also, there realistically should not typically be any inv + // messages with a type of block for the same reason. However, it + // doesn't hurt to update the state accordingly just in case. + // Add the block to the cache of known inventory for the peer. This // helps avoid sending blocks to the peer that it is already known // to have. @@ -1095,22 +1180,6 @@ func (m *SyncManager) handleInvMsg(imsg *invMsg) { // Update the last block in the announced inventory. lastBlock = iv - // Ignore block announcements for blocks that are from peers other - // than the sync peer before the chain is current or are otherwise - // not needed, such as when they're already available. - // - // This helps prevent fetching a mass of orphans. - if (!fromSyncPeer && !isCurrent) || !m.needBlock(&iv.Hash) { - continue - } - - // Request the block if there is not one already pending. - if _, exists := m.requestedBlocks[iv.Hash]; !exists { - limitAdd(m.requestedBlocks, iv.Hash, maxRequestedBlocks) - limitAdd(peer.requestedBlocks, iv.Hash, maxRequestedBlocks) - requestQueue = append(requestQueue, iv) - } - case wire.InvTypeTx: // Add the tx to the cache of known inventory for the peer. This // helps avoid sending transactions to the peer that it is already @@ -1149,15 +1218,6 @@ func (m *SyncManager) handleInvMsg(imsg *invMsg) { peer.UpdateLastBlockHeight(int64(header.Height)) } } - - // Request blocks after the last one advertised up to the final one the - // remote peer knows about (zero stop hash). - blkLocator := m.cfg.Chain.BlockLocatorFromHash(&lastBlock.Hash) - locator := chainBlockLocatorToHashes(blkLocator) - err := peer.PushGetBlocksMsg(locator, &zeroHash) - if err != nil { - log.Errorf("Failed to push getblocksmsg: %v", err) - } } // Request as much as possible at once. @@ -1292,6 +1352,18 @@ out: log.Warnf("Invalid message type in event handler: %T", msg) } + case <-m.hdrSyncState.stallTimer.C: + // Mark the timer's channel as having been drained so the timer can + // safely be reset. + m.hdrSyncState.stallChanDrained = true + + // Disconnect the sync peer due to stalling the header sync process. + if m.syncPeer != nil { + log.Debugf("Header sync progress stalled from peer %s -- "+ + "disconnecting", m.syncPeer) + m.syncPeer.Disconnect() + } + case <-ctx.Done(): break out } @@ -1386,13 +1458,15 @@ func (m *SyncManager) RequestFromPeer(p *peerpkg.Peer, blocks, voteHashes, tSpendHashes []*chainhash.Hash) error { reply := make(chan requestFromPeerResponse, 1) - select { - case m.msgChan <- requestFromPeerMsg{ + request := requestFromPeerMsg{ peer: p, blocks: blocks, voteHashes: voteHashes, tSpendHashes: tSpendHashes, - reply: reply}: + reply: reply, + } + select { + case m.msgChan <- request: case <-m.quit: } @@ -1616,10 +1690,6 @@ type Config struct { // It may return nil if there is no active RPC server. RpcServer func() *rpcserver.Server - // DisableCheckpoints indicates whether or not the sync manager should make - // use of checkpoints. - DisableCheckpoints bool - // NoMiningStateSync indicates whether or not the sync manager should // perform an initial mining state synchronization with peers once they are // believed to be fully synced. @@ -1643,33 +1713,17 @@ type Config struct { // New returns a new network chain synchronization manager. Use Run to begin // processing asynchronous events. func New(config *Config) *SyncManager { - m := SyncManager{ + return &SyncManager{ cfg: *config, rejectedTxns: make(map[chainhash.Hash]struct{}), requestedTxns: make(map[chainhash.Hash]struct{}), requestedBlocks: make(map[chainhash.Hash]struct{}), peers: make(map[*peerpkg.Peer]*syncMgrPeer), + hdrSyncState: makeHeaderSyncState(), progressLogger: progresslog.New("Processed", log), msgChan: make(chan interface{}, config.MaxPeers*3), - headerList: list.New(), quit: make(chan struct{}), + syncHeight: config.Chain.BestSnapshot().Height, isCurrent: config.Chain.IsCurrent(), } - - best := m.cfg.Chain.BestSnapshot() - if !m.cfg.DisableCheckpoints { - // Initialize the next checkpoint based on the current height. - m.nextCheckpoint = m.findNextHeaderCheckpoint(best.Height) - if m.nextCheckpoint != nil { - m.resetHeaderState(&best.Hash, best.Height) - } - } else { - log.Info("Checkpoints are disabled") - } - - m.syncHeightMtx.Lock() - m.syncHeight = best.Height - m.syncHeightMtx.Unlock() - - return &m } diff --git a/peer/log.go b/peer/log.go index 30feeb35..6c7233b7 100644 --- a/peer/log.go +++ b/peer/log.go @@ -1,5 +1,5 @@ // Copyright (c) 2015-2016 The btcsuite developers -// Copyright (c) 2016-2020 The Decred developers +// Copyright (c) 2016-2021 The Decred developers // Use of this source code is governed by an ISC // license that can be found in the LICENSE file. @@ -188,7 +188,13 @@ func messageSummary(msg wire.Message) string { return locatorSummary(msg.BlockLocatorHashes, &msg.HashStop) case *wire.MsgHeaders: - return fmt.Sprintf("num %d", len(msg.Headers)) + summary := fmt.Sprintf("num %d", len(msg.Headers)) + if len(msg.Headers) > 0 { + finalHeader := msg.Headers[len(msg.Headers)-1] + summary = fmt.Sprintf("%s, final hash %s, height %d", summary, + finalHeader.BlockHash(), finalHeader.Height) + } + return summary case *wire.MsgReject: // Ensure the variable length strings don't contain any diff --git a/server.go b/server.go index 6db83815..6bafdad3 100644 --- a/server.go +++ b/server.go @@ -673,7 +673,7 @@ func (sp *serverPeer) OnVersion(p *peer.Peer, msg *wire.MsgVersion) *wire.MsgRej // Ignore peers that have a protocol version that is too old. The peer // negotiation logic will disconnect it after this callback returns. - if msg.ProtocolVersion < int32(wire.InitialProcotolVersion) { + if msg.ProtocolVersion < int32(wire.SendHeadersVersion) { return nil } @@ -737,6 +737,13 @@ func (sp *serverPeer) OnVersion(p *peer.Peer, msg *wire.MsgVersion) *wire.MsgRej return nil } +// OnVerAck is invoked when a peer receives a verack wire message. It creates +// and sends a sendheaders message to request all block annoucements are made +// via full headers instead of the inv message. +func (sp *serverPeer) OnVerAck(_ *peer.Peer, msg *wire.MsgVerAck) { + sp.QueueMessage(wire.NewMsgSendHeaders(), nil) +} + // OnMemPool is invoked when a peer receives a mempool wire message. It creates // and sends an inventory message with the contents of the memory pool up to the // maximum inventory allowed per message. @@ -1067,12 +1074,6 @@ func (sp *serverPeer) OnInv(p *peer.Peer, msg *wire.MsgInv) { // OnHeaders is invoked when a peer receives a headers wire message. The // message is passed down to the net sync manager. func (sp *serverPeer) OnHeaders(_ *peer.Peer, msg *wire.MsgHeaders) { - // Ban peers sending empty headers requests. - if len(msg.Headers) == 0 { - sp.server.BanPeer(sp) - return - } - sp.server.syncManager.QueueHeaders(msg, sp.Peer) } @@ -2068,6 +2069,7 @@ func newPeerConfig(sp *serverPeer) *peer.Config { return &peer.Config{ Listeners: peer.MessageListeners{ OnVersion: sp.OnVersion, + OnVerAck: sp.OnVerAck, OnMemPool: sp.OnMemPool, OnGetMiningState: sp.OnGetMiningState, OnMiningState: sp.OnMiningState, @@ -3508,6 +3510,11 @@ func newServer(ctx context.Context, listenAddrs []string, db database.DB, chainP }, } s.txMemPool = mempool.New(&txC) + + // Create a new sync manager instance with the appropriate configuration. + if cfg.DisableCheckpoints { + srvrLog.Info("Checkpoints are disabled") + } s.syncManager = netsync.New(&netsync.Config{ PeerNotifier: &s, Chain: s.chain, @@ -3516,7 +3523,6 @@ func newServer(ctx context.Context, listenAddrs []string, db database.DB, chainP RpcServer: func() *rpcserver.Server { return s.rpcServer }, - DisableCheckpoints: cfg.DisableCheckpoints, NoMiningStateSync: cfg.NoMiningStateSync, MaxPeers: cfg.MaxPeers, MaxOrphanTxs: cfg.MaxOrphanTxs,