eth: downloader + queue: more context in errors

This commit is contained in:
Martin Holst Swende 2020-05-12 16:39:15 +02:00
parent 0b63915430
commit 236b29aeaf
No known key found for this signature in database
GPG key ID: 683B438C05A5DDF0
2 changed files with 36 additions and 19 deletions

View file

@ -321,13 +321,28 @@ func (d *Downloader) UnregisterPeer(id string) error {
// adding various sanity checks as well as wrapping it with various log entries. // adding various sanity checks as well as wrapping it with various log entries.
func (d *Downloader) Synchronise(id string, head common.Hash, td *big.Int, mode SyncMode) error { func (d *Downloader) Synchronise(id string, head common.Hash, td *big.Int, mode SyncMode) error {
err := d.synchronise(id, head, td, mode) err := d.synchronise(id, head, td, mode)
switch err {
case nil:
case errBusy, errCanceled:
switch err {
case nil, errBusy, errCanceled:
return err
}
if errors.Is(err, errInvalidChain) {
log.Warn("Synchronisation failed, dropping peer", "peer", id, "err", err)
if d.dropPeer == nil {
// The dropPeer method is nil when `--copydb` is used for a local copy.
// Timeouts can occur if e.g. compaction hits at the wrong time, and can be ignored
log.Warn("Downloader wants to drop peer, but peerdrop-function is not set", "peer", id)
} else {
d.dropPeer(id)
}
return err
}
switch err {
case errTimeout, errBadPeer, errStallingPeer, errUnsyncedPeer, case errTimeout, errBadPeer, errStallingPeer, errUnsyncedPeer,
errEmptyHeaderSet, errPeersUnavailable, errTooOld, errEmptyHeaderSet, errPeersUnavailable, errTooOld,
errInvalidAncestor, errInvalidChain: errInvalidAncestor:
log.Warn("Synchronisation failed, dropping peer", "peer", id, "err", err) log.Warn("Synchronisation failed, dropping peer", "peer", id, "err", err)
if d.dropPeer == nil { if d.dropPeer == nil {
// The dropPeer method is nil when `--copydb` is used for a local copy. // The dropPeer method is nil when `--copydb` is used for a local copy.
@ -774,7 +789,7 @@ func (d *Downloader) findAncestor(p *peerConnection, remoteHeader *types.Header)
expectNumber := from + int64(i)*int64(skip+1) expectNumber := from + int64(i)*int64(skip+1)
if number := header.Number.Int64(); number != expectNumber { if number := header.Number.Int64(); number != expectNumber {
p.log.Warn("Head headers broke chain ordering", "index", i, "requested", expectNumber, "received", number) p.log.Warn("Head headers broke chain ordering", "index", i, "requested", expectNumber, "received", number)
return 0, errInvalidChain return 0, fmt.Errorf("%w: %v", errInvalidChain, errors.New("head headers broke chain ordering"))
} }
} }
// Check if a common ancestor was found // Check if a common ancestor was found
@ -988,7 +1003,7 @@ func (d *Downloader) fetchHeaders(p *peerConnection, from uint64, pivot uint64)
filled, proced, err := d.fillHeaderSkeleton(from, headers) filled, proced, err := d.fillHeaderSkeleton(from, headers)
if err != nil { if err != nil {
p.log.Debug("Skeleton chain invalid", "err", err) p.log.Debug("Skeleton chain invalid", "err", err)
return errInvalidChain return fmt.Errorf("%w: %v", errInvalidChain, err)
} }
headers = filled[proced:] headers = filled[proced:]
from += uint64(proced) from += uint64(proced)
@ -1207,13 +1222,13 @@ func (d *Downloader) fetchParts(deliveryCh chan dataPack, deliver func(dataPack)
if peer := d.peers.Peer(packet.PeerId()); peer != nil { if peer := d.peers.Peer(packet.PeerId()); peer != nil {
// Deliver the received chunk of data and check chain validity // Deliver the received chunk of data and check chain validity
accepted, err := deliver(packet) accepted, err := deliver(packet)
if err == errInvalidChain { if errors.Is(err, errInvalidChain) {
return err return err
} }
// Unless a peer delivered something completely else than requested (usually // Unless a peer delivered something completely else than requested (usually
// caused by a timed out request which came through in the end), set it to // caused by a timed out request which came through in the end), set it to
// idle. If the delivery's stale, the peer should have already been idled. // idle. If the delivery's stale, the peer should have already been idled.
if err != errStaleDelivery { if !errors.Is(err, errStaleDelivery) {
setIdle(peer, accepted) setIdle(peer, accepted)
} }
// Issue a log to the user to see what's going on // Issue a log to the user to see what's going on
@ -1473,7 +1488,7 @@ func (d *Downloader) processHeaders(origin uint64, pivot uint64, td *big.Int) er
rollback = append(rollback, chunk[:n]...) rollback = append(rollback, chunk[:n]...)
} }
log.Debug("Invalid header encountered", "number", chunk[n].Number, "hash", chunk[n].Hash(), "err", err) log.Debug("Invalid header encountered", "number", chunk[n].Number, "hash", chunk[n].Hash(), "err", err)
return errInvalidChain return fmt.Errorf("%w: %v", errInvalidChain, err)
} }
// All verifications passed, store newly found uncertain headers // All verifications passed, store newly found uncertain headers
rollback = append(rollback, unknown...) rollback = append(rollback, unknown...)
@ -1565,7 +1580,7 @@ func (d *Downloader) importBlockResults(results []*fetchResult) error {
// of the blocks delivered from the downloader, and the indexing will be off. // of the blocks delivered from the downloader, and the indexing will be off.
log.Debug("Downloaded item processing failed on sidechain import", "index", index, "err", err) log.Debug("Downloaded item processing failed on sidechain import", "index", index, "err", err)
} }
return errInvalidChain return fmt.Errorf("%w: %v", errInvalidChain, err)
} }
return nil return nil
} }
@ -1706,7 +1721,7 @@ func (d *Downloader) commitFastSyncData(results []*fetchResult, stateSync *state
} }
if index, err := d.blockchain.InsertReceiptChain(blocks, receipts, d.ancientLimit); err != nil { if index, err := d.blockchain.InsertReceiptChain(blocks, receipts, d.ancientLimit); err != nil {
log.Debug("Downloaded item processing failed", "number", results[index].Header.Number, "hash", results[index].Header.Hash(), "err", err) log.Debug("Downloaded item processing failed", "number", results[index].Header.Number, "hash", results[index].Header.Hash(), "err", err)
return errInvalidChain return fmt.Errorf("%w: %v", errInvalidChain, err)
} }
return nil return nil
} }

View file

@ -509,7 +509,7 @@ func (q *queue) reserveHeaders(p *peerConnection, count int, taskPool map[common
index := int(header.Number.Int64() - int64(q.resultOffset)) index := int(header.Number.Int64() - int64(q.resultOffset))
if index >= len(q.resultCache) || index < 0 { if index >= len(q.resultCache) || index < 0 {
common.Report("index allocation went beyond available resultCache space") common.Report("index allocation went beyond available resultCache space")
return nil, false, errInvalidChain return nil, false, fmt.Errorf("%w: index allocation went beyond available resultCache space", errInvalidChain)
} }
if q.resultCache[index] == nil { if q.resultCache[index] == nil {
components := 1 components := 1
@ -863,14 +863,16 @@ func (q *queue) deliver(id string, taskPool map[common.Hash]*types.Header, taskQ
q.active.Signal() q.active.Signal()
} }
// If none of the data was good, it's a stale delivery // If none of the data was good, it's a stale delivery
switch { if failure == nil {
case failure == nil || failure == errInvalidChain: return accepted, nil
return accepted, failure
case useful:
return accepted, fmt.Errorf("partial failure: %v", failure)
default:
return accepted, errStaleDelivery
} }
if errors.Is(failure, errInvalidChain) {
return accepted, failure
}
if useful {
return accepted, fmt.Errorf("partial failure: %v", failure)
}
return accepted, fmt.Errorf("%w: %v", failure, errStaleDelivery)
} }
// Prepare configures the result cache to allow accepting and caching inbound // Prepare configures the result cache to allow accepting and caching inbound