From 373f600c7b94cd515d7af8691986f192211a6be1 Mon Sep 17 00:00:00 2001 From: Martin Holst Swende Date: Sat, 27 Jan 2018 12:30:28 +0100 Subject: [PATCH] core, downloader, handler, sync: make use of mutex-free stats --- common/chainstats/chainstats.go | 4 ++++ core/blockchain.go | 18 ++++++++++++++++++ eth/downloader/downloader.go | 28 ++++++++++++++++++---------- eth/handler.go | 8 ++++---- eth/sync.go | 10 +++++----- 5 files changed, 49 insertions(+), 19 deletions(-) diff --git a/common/chainstats/chainstats.go b/common/chainstats/chainstats.go index 35397eff66..7df0a3f884 100644 --- a/common/chainstats/chainstats.go +++ b/common/chainstats/chainstats.go @@ -47,6 +47,10 @@ func (stats *Chainstats) SetNumber(number *big.Int) { func (stats *Chainstats) GetFastNumber() uint64 { return stats.currentFastBlockNumber.Load().(*big.Int).Uint64() } +func (stats *Chainstats) GetNumbers() (uint64, uint64) { + return stats.currentBlockNumber.Load().(*big.Int).Uint64(), + stats.currentFastBlockNumber.Load().(*big.Int).Uint64() +} func (stats *Chainstats) SetFastNumber(number *big.Int) { stats.currentFastBlockNumber.Store(number) } diff --git a/core/blockchain.go b/core/blockchain.go index 9adecef7a9..c7ce127368 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -714,6 +714,24 @@ func (bc *BlockChain) procFutureBlocks() { } } +// Stats returns the chainstats which can be queried non-blocking for info about difficulty and numbers +// These are provided on a best-effort, and it's theoretically possible that two consecutive calls to +// number and difficulty return number for X and difficulty for Y, if the stats is updated between the calls +func (bc *BlockChain) Stats() *chainstats.Chainstats { + return bc.chainStats +} +func (bc *BlockChain) CurrentNumber() uint64 { + return bc.chainStats.GetNumber() +} + +func (bc *BlockChain) CurrentFastNumber() uint64 { + return bc.chainStats.GetFastNumber() +} + +func (bc *BlockChain) CurrentTD() *big.Int { + return bc.chainStats.GetTotalDifficulty() +} + // WriteStatus status of write type WriteStatus byte diff --git a/eth/downloader/downloader.go b/eth/downloader/downloader.go index 7ede530a94..32cbca6051 100644 --- a/eth/downloader/downloader.go +++ b/eth/downloader/downloader.go @@ -193,6 +193,15 @@ type BlockChain interface { // InsertReceiptChain inserts a batch of receipts into the local chain. InsertReceiptChain(types.Blocks, []types.Receipts) (int, error) + + // CurrentNumber retrieves number of the current head block + CurrentNumber() uint64 + + // CurrentNumber retrieves number of the current head fast block + CurrentFastNumber() uint64 + + // CurrentTD retrives the total difficulty of the current head block + CurrentTD() *big.Int } // New creates a new downloader to fetch hashes and blocks from remote peers. @@ -583,9 +592,9 @@ func (d *Downloader) findAncestor(p *peerConnection, height uint64) (uint64, err floor, ceil := int64(-1), d.lightchain.CurrentHeader().Number.Uint64() if d.mode == FullSync { - ceil = d.blockchain.CurrentBlock().NumberU64() + ceil = d.blockchain.CurrentNumber() } else if d.mode == FastSync { - ceil = d.blockchain.CurrentFastBlock().NumberU64() + ceil = d.blockchain.CurrentFastNumber() } if ceil >= MaxForkAncestry { floor = int64(ceil - MaxForkAncestry) @@ -1156,16 +1165,16 @@ func (d *Downloader) processHeaders(origin uint64, pivot uint64, td *big.Int) er for i, header := range rollback { hashes[i] = header.Hash() } - lastHeader, lastFastBlock, lastBlock := d.lightchain.CurrentHeader().Number, common.Big0, common.Big0 + lastHeader, lastFastBlock, lastBlock := d.lightchain.CurrentHeader().Number, uint64(0), uint64(0) if d.mode != LightSync { - lastFastBlock = d.blockchain.CurrentFastBlock().Number() - lastBlock = d.blockchain.CurrentBlock().Number() + lastFastBlock = d.blockchain.CurrentFastNumber() + lastBlock = d.blockchain.CurrentNumber() } d.lightchain.Rollback(hashes) - curFastBlock, curBlock := common.Big0, common.Big0 + curFastBlock, curBlock := uint64(0), uint64(0) if d.mode != LightSync { - curFastBlock = d.blockchain.CurrentFastBlock().Number() - curBlock = d.blockchain.CurrentBlock().Number() + curFastBlock = d.blockchain.CurrentFastNumber() + curBlock = d.blockchain.CurrentNumber() } log.Warn("Rolled back headers", "count", len(hashes), "header", fmt.Sprintf("%d->%d", lastHeader, d.lightchain.CurrentHeader().Number), @@ -1205,8 +1214,7 @@ func (d *Downloader) processHeaders(origin uint64, pivot uint64, td *big.Int) er // L: Request new headers up from 11 (R's TD was higher, it must have something) // R: Nothing to give if d.mode != LightSync { - head := d.blockchain.CurrentBlock() - if !gotHeaders && td.Cmp(d.blockchain.GetTd(head.Hash(), head.NumberU64())) > 0 { + if !gotHeaders && td.Cmp(d.blockchain.CurrentTD()) > 0 { return errStallingPeer } } diff --git a/eth/handler.go b/eth/handler.go index c2426544f6..71a0ee36ee 100644 --- a/eth/handler.go +++ b/eth/handler.go @@ -113,7 +113,7 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne quitSync: make(chan struct{}), } // Figure out whether to allow fast sync or not - if mode == downloader.FastSync && blockchain.CurrentBlock().NumberU64() > 0 { + if mode == downloader.FastSync && blockchain.Stats().GetNumber() > 0 { log.Warn("Blockchain not empty, fast sync disabled") mode = downloader.FullSync } @@ -165,7 +165,7 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne return engine.VerifyHeader(blockchain, header, true) } heighter := func() uint64 { - return blockchain.CurrentBlock().NumberU64() + return blockchain.Stats().GetNumber() } inserter := func(blocks types.Blocks) (int, error) { // If fast sync is running, deny importing weird blocks @@ -647,8 +647,8 @@ func (pm *ProtocolManager) handleMsg(p *peer) error { // Schedule a sync if above ours. Note, this will not fire a sync for a gap of // a singe block (as the true TD is below the propagated block), however this // scenario should easily be covered by the fetcher. - currentBlock := pm.blockchain.CurrentBlock() - if trueTD.Cmp(pm.blockchain.GetTd(currentBlock.Hash(), currentBlock.NumberU64())) > 0 { + currentTd := pm.blockchain.Stats().GetTotalDifficulty() + if trueTD.Cmp(currentTd) > 0 { go pm.synchronise(p) } } diff --git a/eth/sync.go b/eth/sync.go index 2da1464bc5..f20fbe55c4 100644 --- a/eth/sync.go +++ b/eth/sync.go @@ -167,8 +167,8 @@ func (pm *ProtocolManager) synchronise(peer *peer) { return } // Make sure the peer's TD is higher than our own - currentBlock := pm.blockchain.CurrentBlock() - td := pm.blockchain.GetTd(currentBlock.Hash(), currentBlock.NumberU64()) + currentNumber, currentFastNumber := pm.blockchain.Stats().GetNumbers() + td := pm.blockchain.Stats().GetTotalDifficulty() pHead, pTd := peer.Head() if pTd.Cmp(td) <= 0 { @@ -179,7 +179,7 @@ func (pm *ProtocolManager) synchronise(peer *peer) { if atomic.LoadUint32(&pm.fastSync) == 1 { // Fast sync was explicitly requested, and explicitly granted mode = downloader.FastSync - } else if currentBlock.NumberU64() == 0 && pm.blockchain.CurrentFastBlock().NumberU64() > 0 { + } else if currentNumber == 0 && currentFastNumber > 0 { // The database seems empty as the current block is the genesis. Yet the fast // block is ahead, so fast sync was enabled for this node at a certain point. // The only scenario where this can happen is if the user manually (or via a @@ -197,13 +197,13 @@ func (pm *ProtocolManager) synchronise(peer *peer) { atomic.StoreUint32(&pm.fastSync, 0) } atomic.StoreUint32(&pm.acceptTxs, 1) // Mark initial sync done - if head := pm.blockchain.CurrentBlock(); head.NumberU64() > 0 { + if pm.blockchain.Stats().GetNumber() > 0 { // We've completed a sync cycle, notify all peers of new state. This path is // essential in star-topology networks where a gateway node needs to notify // all its out-of-date peers of the availability of a new block. This failure // scenario will most often crop up in private and hackathon networks with // degenerate connectivity, but it should be healthy for the mainnet too to // more reliably update peers or the local TD state. - go pm.BroadcastBlock(head, false) + go pm.BroadcastBlock(pm.blockchain.CurrentBlock(), false) } }