Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
default: major
---

# Compared header chain work before downloading blocks from a peer

The syncer now walks a peer's headers until their chain either outweighs ours or runs out, and only downloads blocks once it knows the chain is worth adopting. Previously a peer stuck on a fork that would never outweigh our chain had its entire fork re-downloaded every sync interval.
10 changes: 5 additions & 5 deletions syncer/peer.go
Original file line number Diff line number Diff line change
Expand Up @@ -140,20 +140,20 @@ func (p *Peer) DiscoverIP(timeout time.Duration) (string, error) {
}

// SendHeaders requests up to n headers from p, starting from the supplied
// index, which must be on the peer's best chain. The peer also returns the
// number of remaining headers left to sync.
func (p *Peer) SendHeaders(cs consensus.State, maxHeaders uint64, timeout time.Duration) ([]types.BlockHeader, uint64, error) {
// index, which must be on the peer's best chain. It also returns the state
// after applying them and the number of headers the peer still has.
func (p *Peer) SendHeaders(cs consensus.State, maxHeaders uint64, timeout time.Duration) ([]types.BlockHeader, consensus.State, uint64, error) {
r := &gateway.RPCSendHeaders{Index: cs.Index, Max: maxHeaders}
err := p.callRPC(r, timeout)
if err == nil {
for _, bh := range r.Headers {
if err := consensus.ValidateHeader(cs, bh); err != nil {
return nil, 0, fmt.Errorf("peer sent invalid header %v: %w", bh.ID(), err)
return nil, consensus.State{}, 0, fmt.Errorf("peer sent invalid header %v: %w", bh.ID(), err)
}
cs = consensus.ApplyHeader(cs, bh, time.Time{})
Comment thread
ChrisSchinnerl marked this conversation as resolved.
}
}
return r.Headers, r.Remaining, err
return r.Headers, cs, r.Remaining, err
}

// SendTransactions requests a subset of a block's transactions from the peer.
Expand Down
205 changes: 169 additions & 36 deletions syncer/syncer.go
Original file line number Diff line number Diff line change
Expand Up @@ -781,6 +781,154 @@ func (s *Syncer) peerLoop(ctx context.Context) error {
return nil
}

// a headerBatch is one SendHeaders response and the state it extends.
type headerBatch struct {
cs consensus.State
headers []types.BlockHeader
}

// a peerChain is the result of walking a peer's header chain.
type peerChain struct {
fork consensus.State // state at the fork point
tip types.BlockHeader // last header needed to outweigh our chain
batchTips []types.BlockID // last header of each batch, to detect a reorg
remaining uint64 // headers the peer has beyond tip
heavier bool
retained []headerBatch // the walked batches, if kept
}

// batches is the number of SendHeaders calls needed to reach tip.
func (pc peerChain) batches() int { return len(pc.batchTips) }

// headersRetained reports whether every walked batch was kept.
func (pc peerChain) headersRetained() bool { return len(pc.retained) == len(pc.batchTips) }

const (
// retainedBatches bounds how many batches a walk keeps
retainedBatches = 2

// maxHeaderWalk bounds one peer's walk; a sync round waits for every peer,
// and a peer that cannot finish within it is dropped
maxHeaderWalk = 2 * time.Minute

// maxChainSync backstops a wedged download; a sync round waits for every
// peer, and parallelSync only bounds individual requests
maxChainSync = 15 * time.Minute
)

var (
// errWalkAbandoned wraps the reasons a walk stops without a verdict; the
// peer stays eligible and is retried next tick
errWalkAbandoned = errors.New("header walk abandoned")
errPeerReorg = fmt.Errorf("%w: peer reorged", errWalkAbandoned)

// errPeerHeaders marks a failure attributable to the peer, unlike a block
// download, which may have involved any peer
errPeerHeaders = errors.New("peer failed to serve headers")
errWalkTimeout = fmt.Errorf("%w: walk exceeded %v", errPeerHeaders, maxHeaderWalk)
)

// walkPeerChain requests headers from p until its chain outweighs ts or it runs
// out of them, reporting whether the chain is worth downloading blocks for.
func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.State) (peerChain, error) {
Comment thread
ChrisSchinnerl marked this conversation as resolved.
// the deadline covers history probing too, and clamps each request
deadline := time.Now().Add(maxHeaderWalk)
sendHeaders := func(cs consensus.State) ([]types.BlockHeader, consensus.State, uint64, error) {
timeout := min(s.config.SendHeadersTimeout, time.Until(deadline))
if timeout <= 0 {
return nil, consensus.State{}, 0, errWalkTimeout
}
headers, tip, remaining, err := p.SendHeaders(cs, s.config.MaxSendHeaders, timeout)
if err != nil && !time.Now().Before(deadline) {
return nil, consensus.State{}, 0, errWalkTimeout
}
return headers, tip, remaining, err
}
for _, id := range hist {
if id == (types.BlockID{}) {
// skip empty history entries which can occur when we don't have a
// full history of blocks.
continue
}
fork, ok := s.cm.State(id)
if !ok {
return peerChain{}, errors.New("missing state for history")
}
headers, tip, remaining, err := sendHeaders(fork)
if err != nil && strings.Contains(err.Error(), "EOF") {
continue // probably "index is not on our best chain"
} else if err != nil {
return peerChain{}, err
}
pc, cs := peerChain{fork: fork}, fork
for {
if len(headers) > 0 {
pc.tip = headers[len(headers)-1]
pc.batchTips = append(pc.batchTips, pc.tip.ID())
if len(pc.batchTips) <= retainedBatches {
pc.retained = append(pc.retained, headerBatch{cs: cs, headers: headers})
} else {
pc.retained = nil
}
}
Comment thread
ChrisSchinnerl marked this conversation as resolved.
pc.remaining = remaining
if tip.SufficientlyHeavierThan(ts) {
pc.heavier = true
return pc, nil
} else if remaining == 0 || len(headers) == 0 {
return pc, nil
}
cs = tip
if headers, tip, remaining, err = sendHeaders(cs); err != nil {
if strings.Contains(err.Error(), "EOF") {
return peerChain{}, errPeerReorg
}
return peerChain{}, err
}
}
}
return peerChain{}, errors.New("no common history")
}

// syncChain downloads the blocks of the chain pc describes, returning how much
// was actually synced. p only serves headers the walk did not retain; blocks
// come from every eligible peer via parallelSync.
func (s *Syncer) syncChain(ctx context.Context, p *Peer, pc peerChain) (peerChain, error) {
ctx, cancel := context.WithTimeout(ctx, maxChainSync)
defer cancel()

synced := peerChain{fork: pc.fork}
if pc.headersRetained() {
for _, b := range pc.retained {
if err := s.parallelSync(ctx, b.cs, b.headers); err != nil {
return synced, err
}
synced.tip = b.headers[len(b.headers)-1]
synced.batchTips = append(synced.batchTips, synced.tip.ID())
}
synced.remaining = pc.remaining
return synced, nil
}
cs := pc.fork
for _, want := range pc.batchTips {
headers, tip, remaining, err := p.SendHeaders(cs, s.config.MaxSendHeaders, s.config.SendHeadersTimeout)
if err != nil && strings.Contains(err.Error(), "EOF") {
return synced, errPeerReorg
} else if err != nil {
return synced, fmt.Errorf("%w: %w", errPeerHeaders, err)
} else if len(headers) == 0 || headers[len(headers)-1].ID() != want {
// they are no longer serving the chain we judged heavier
return synced, errPeerReorg
} else if err := s.parallelSync(ctx, cs, headers); err != nil {
Comment thread
ChrisSchinnerl marked this conversation as resolved.
return synced, err
}
cs = tip
synced.tip, synced.remaining = headers[len(headers)-1], remaining
synced.batchTips = append(synced.batchTips, synced.tip.ID())
}
return synced, nil
}

func (s *Syncer) syncLoop(ctx context.Context) error {
ticker := time.NewTicker(s.config.SyncInterval)
defer ticker.Stop()
Expand All @@ -801,62 +949,47 @@ func (s *Syncer) syncLoop(ctx context.Context) error {
}
s.mu.Unlock()
type resp struct {
peer *Peer
cs consensus.State
headers []types.BlockHeader
remaining uint64
err error
peer *Peer
chain peerChain
err error
}
respChan := make(chan resp, len(peers))
hist, err := s.cm.History()
if err != nil {
return err // generally fatal
}
tipState := s.cm.TipState()
for _, p := range peers {
go func(p *Peer) {
cs, headers, remaining, err := func() (consensus.State, []types.BlockHeader, uint64, error) {
for _, id := range hist {
if id == (types.BlockID{}) {
// skip empty history entries which can occur when
// we don't have a full history of blocks.
continue
}
cs, ok := s.cm.State(id)
if !ok {
return consensus.State{}, nil, 0, errors.New("missing state for history")
}
headers, remaining, err := p.SendHeaders(cs, s.config.MaxSendHeaders, s.config.SendHeadersTimeout)
if err != nil && strings.Contains(err.Error(), "EOF") {
continue // probably "index is not on our best chain"
} else if err != nil {
return consensus.State{}, nil, 0, err
}
return cs, headers, remaining, nil
}
return consensus.State{}, nil, 0, errors.New("no common history")
}()
respChan <- resp{peer: p, cs: cs, headers: headers, remaining: remaining, err: err}
pc, err := s.walkPeerChain(p, hist, tipState)
respChan <- resp{peer: p, chain: pc, err: err}
}(p)
}
// sync each set of headers as they arrive
seen := make(map[types.BlockID]bool)
for range peers {
if r := <-respChan; r.err != nil {
if r := <-respChan; errors.Is(r.err, errWalkAbandoned) {
Comment thread
ChrisSchinnerl marked this conversation as resolved.
s.log.Debug("abandoned header walk", zap.Stringer("peer", r.peer), zap.Error(r.err))
} else if r.err != nil {
r.peer.setErr(r.err)
} else if len(r.headers) == 0 {
} else if !r.chain.heavier {
// a later relay from them triggers a resync
s.log.Debug("peer chain does not outweigh ours", zap.Stringer("peer", r.peer), zap.Int("batches", r.chain.batches()))
r.peer.setSynced(true)
} else if id := r.headers[len(r.headers)-1].ID(); seen[id] {
} else if id := r.chain.tip.ID(); seen[id] {
continue // already syncing these blocks from another peer
} else {
seen[id] = true
s.log.Debug("syncing blocks", zap.Stringer("peer", r.peer), zap.Stringer("start", r.cs.Index), zap.Int("n", len(r.headers)))
if err := s.parallelSync(ctx, r.cs, r.headers); err != nil {
s.log.Debug("syncing blocks", zap.Stringer("peer", r.peer), zap.Stringer("start", r.chain.fork.Index), zap.Int("batches", r.chain.batches()))
synced, err := s.syncChain(ctx, r.peer, r.chain)
if errors.Is(err, errPeerHeaders) {
r.peer.setErr(err)
} else if err != nil {
s.log.Debug("sync failed", zap.Stringer("peer", r.peer), zap.Error(err))
} else if r.remaining == 0 {
// peer sent all their headers; mark them as synced and
// relay their tip
} else if synced.batches() > 0 && synced.remaining == 0 {
// peer sent all their headers
r.peer.setSynced(true)
go s.relayV2Header(r.headers[len(r.headers)-1], r.peer)
go s.relayV2Header(synced.tip, r.peer)
}
}
}
Expand Down
Loading
Loading