diff --git a/eth/fetcher/fetcher.go b/eth/fetcher/fetcher.go index 94f05f9674..8a3aa2cb83 100644 --- a/eth/fetcher/fetcher.go +++ b/eth/fetcher/fetcher.go @@ -67,6 +67,12 @@ type chainInsertFn func(types.Blocks) (int, error) // peerDropFn is a callback type for dropping a peer detected as malicious. type peerDropFn func(id string) +// peerHeaderRequesterFn is a callback type for sending a header retrieval request. +type peerHeaderRequesterFn func(string) func(common.Hash) error + +// peerBodyRequesterFn is a callback type for sending a body retrieval request. +type peerBodyRequesterFn func(string) func([]common.Hash) error + // announce is the hash notification of the availability of a new block in the // network. type announce struct { @@ -135,6 +141,8 @@ type Fetcher struct { chainHeight chainHeightFn // Retrieves the current chain's height insertChain chainInsertFn // Injects a batch of blocks into the chain dropPeer peerDropFn // Drops a peer for misbehaving + headerRequest peerHeaderRequesterFn + bodyRequest peerBodyRequesterFn // Testing hooks announceChangeHook func(common.Hash, bool) // Method to call upon adding or deleting a hash from the announce list @@ -145,7 +153,7 @@ type Fetcher struct { } // New creates a block fetcher to retrieve blocks based on hash announcements. -func New(getBlock blockRetrievalFn, verifyHeader headerVerifierFn, broadcastBlock blockBroadcasterFn, chainHeight chainHeightFn, insertChain chainInsertFn, dropPeer peerDropFn) *Fetcher { +func New(getBlock blockRetrievalFn, verifyHeader headerVerifierFn, broadcastBlock blockBroadcasterFn, chainHeight chainHeightFn, insertChain chainInsertFn, dropPeer peerDropFn, headerRequest peerHeaderRequesterFn, bodyRequest peerBodyRequesterFn) *Fetcher { return &Fetcher{ notify: make(chan *announce), inject: make(chan *inject), @@ -167,6 +175,8 @@ func New(getBlock blockRetrievalFn, verifyHeader headerVerifierFn, broadcastBloc chainHeight: chainHeight, insertChain: insertChain, dropPeer: dropPeer, + headerRequest: headerRequest, + bodyRequest: bodyRequest, } } @@ -296,18 +306,36 @@ func (f *Fetcher) loop() { } // If too high up the chain or phase, continue later number := op.block.NumberU64() - if number > height+1 { - f.queue.Push(op, -int64(number)) - if f.queueChangeHook != nil { - f.queueChangeHook(hash, true) - } - break - } // Otherwise if fresh and still unknown, try and import if number+maxUncleDist < height || f.getBlock(hash) != nil { f.forgetBlock(hash) continue } + parentHash := op.block.ParentHash() + if f.getBlock(parentHash) == nil { + f.queue.Push(op, -int64(number)) + if f.queueChangeHook != nil { + f.queueChangeHook(hash, true) + } + if len(f.announced[parentHash]) > 0 { + // only request for parent hash once + break + } + headerRequestFn := f.headerRequest(op.origin) + bodyRequestFn := f.bodyRequest(op.origin) + if headerRequestFn == nil || bodyRequestFn == nil { + log.Info("Peer lost while fetching for parent", "peer", op.origin) + break + } + go func() { + log.Info("Notify of parent hash", "hash", parentHash, "number", op.block.NumberU64()-1, "peer", op.origin) + err := f.Notify(op.origin, parentHash, op.block.NumberU64()-1, time.Now(), headerRequestFn, bodyRequestFn) + if err != nil { + log.Error("Unable to notify of parent block", "error", err) + } + }() + break + } f.insert(op.origin, op.block) } // Wait for an outside event to occur diff --git a/eth/handler.go b/eth/handler.go index b42612a566..7d286498cf 100644 --- a/eth/handler.go +++ b/eth/handler.go @@ -182,7 +182,19 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne atomic.StoreUint32(&manager.acceptTxs, 1) // Mark initial sync done on any fetcher import return manager.blockchain.InsertChain(blocks) } - manager.fetcher = fetcher.New(blockchain.GetBlockByHash, validator, manager.BroadcastBlock, heighter, inserter, manager.removePeer) + headerRequest := func(id string) func(common.Hash) error { + if peer := manager.peers.Peer(id); peer != nil { + return peer.RequestOneHeader + } + return nil + } + bodyRequest := func(id string) func([]common.Hash) error { + if peer := manager.peers.Peer(id); peer != nil { + return peer.RequestBodies + } + return nil + } + manager.fetcher = fetcher.New(blockchain.GetBlockByHash, validator, manager.BroadcastBlock, heighter, inserter, manager.removePeer, headerRequest, bodyRequest) return manager, nil }