From 93a87b69b810ab7acb38bd3414ad1f400eeff8c6 Mon Sep 17 00:00:00 2001 From: Chris Schinnerl <3903476+ChrisSchinnerl@users.noreply.github.com> Date: Tue, 22 Sep 2026 11:06:48 +0200 Subject: [PATCH 1/5] Stop repeatedly syncing from peers that don't extend our best chain --- ...ose_blocks_do_not_extend_our_best_chain.md | 7 + syncer/syncer.go | 11 ++ syncer/syncer_test.go | 159 ++++++++++++++++++ 3 files changed, 177 insertions(+) create mode 100644 .changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md diff --git a/.changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md b/.changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md new file mode 100644 index 00000000..137e6d94 --- /dev/null +++ b/.changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md @@ -0,0 +1,7 @@ +--- +default: patch +--- + +# Stopped repeatedly syncing from peers whose blocks do not extend our best chain + +The syncer now skips a peer if our tip is unchanged since the last batch it served. Previously a peer stuck on a fork that would never outweigh our chain was queried every sync interval, re-downloading its entire fork indefinitely. diff --git a/syncer/syncer.go b/syncer/syncer.go index ec854d22..ee445a82 100644 --- a/syncer/syncer.go +++ b/syncer/syncer.go @@ -850,8 +850,19 @@ func (s *Syncer) syncLoop(ctx context.Context) error { } 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))) + tip := s.cm.Tip() if err := s.parallelSync(ctx, r.cs, r.headers); err != nil { s.log.Debug("sync failed", zap.Stringer("peer", r.peer), zap.Error(err)) + } else if s.cm.Tip() == tip { + // their blocks did not extend our best chain, so the peer + // has nothing for us. Mark them synced; if they later learn + // of a block we don't have, relaying it triggers a resync. + // + // NOTE: a peer stuck on a chain we will never adopt never + // reports remaining == 0, so without this we would ask them + // for the same headers every sync interval, indefinitely. + s.log.Debug("peer has no blocks that extend our chain", zap.Stringer("peer", r.peer), zap.Stringer("tip", tip)) + r.peer.setSynced(true) } else if r.remaining == 0 { // peer sent all their headers; mark them as synced and // relay their tip diff --git a/syncer/syncer_test.go b/syncer/syncer_test.go index af52d605..72d0bf65 100644 --- a/syncer/syncer_test.go +++ b/syncer/syncer_test.go @@ -862,3 +862,162 @@ func TestMaxInflightRPCsBackpressureNotDropped(t *testing.T) { t.Fatalf("backpressured RPC failed after the slot was freed: %v", err) } } + +// forkManager serves a limited number of headers per request and counts block +// requests, simulating a peer on a fork longer than a single sync batch. +type forkManager struct { + headerLimit uint64 + headerReqs atomic.Uint64 + blockReqs atomic.Uint64 + *chain.Manager +} + +func (fm *forkManager) Headers(index types.ChainIndex, maxHeaders uint64) ([]types.BlockHeader, uint64, error) { + fm.headerReqs.Add(1) + return fm.Manager.Headers(index, min(maxHeaders, fm.headerLimit)) +} + +func (fm *forkManager) BlocksForHistory(history []types.BlockID, maxBlocks uint64) ([]types.Block, uint64, error) { + fm.blockReqs.Add(1) + return fm.Manager.BlocksForHistory(history, maxBlocks) +} + +// TestForkPeerNotResynced verifies that we stop syncing from a peer on a fork +// that never becomes our best chain. Such a peer is never marked synced, +// because it always has headers remaining; since our tip doesn't change, it +// would answer every subsequent sync round with the same headers, and we would +// re-download the same blocks indefinitely. +func TestForkPeerNotResynced(t *testing.T) { + log := zaptest.NewLogger(t) + + // s1 has the longer chain + s1, cm1 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer1"))) + defer s1.Close() + testutil.MineBlocks(t, cm1, types.VoidAddress, 20) + + // s2 is on a shorter fork, and only serves 5 headers at a time, so it + // always reports headers remaining + n, genesis := testutil.Network() + store2, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) + if err != nil { + t.Fatal(err) + } + fm := &forkManager{headerLimit: 5, Manager: chain.NewManager(store2)} + // mine to a different address so the fork diverges at the genesis child, + // leaving more headers than a single batch can carry + testutil.MineBlocks(t, fm.Manager, types.Address{1}, 10) + + l2, err := net.Listen("tcp", ":0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { l2.Close() }) + + s2 := syncer.New(l2, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ + GenesisID: genesis.ID(), + UniqueID: gateway.GenerateUniqueID(), + NetAddress: l2.Addr().String(), + }, syncer.WithSyncInterval(time.Hour)) // effectively disabled + go s2.Run() + defer s2.Close() + + s1Tip := cm1.Tip() + + if _, err := s1.Connect(context.Background(), s2.Addr()); err != nil { + t.Fatal(err) + } + + // wait for s1 to fetch a batch of fork blocks + for range 100 { + if fm.blockReqs.Load() > 0 { + break + } + time.Sleep(100 * time.Millisecond) + } + if fm.blockReqs.Load() == 0 { + t.Fatal("s1 never synced any blocks from the fork peer") + } + + // s1 should give up on the peer rather than starting a fresh sync round + // every sync interval; allow a few in-flight rounds to settle + time.Sleep(time.Second) + settled, settledHeaders := fm.blockReqs.Load(), fm.headerReqs.Load() + time.Sleep(2 * time.Second) // ~20 sync intervals + if reqs := fm.headerReqs.Load(); reqs != settledHeaders { + t.Fatalf("s1 kept starting sync rounds with the fork peer: %v more header requests", reqs-settledHeaders) + } else if reqs := fm.blockReqs.Load(); reqs != settled { + t.Fatalf("s1 kept re-syncing the fork peer: %v block requests, then %v more", settled, reqs-settled) + } + + // the peer had nothing that extends our chain, so it is done with + peers := s1.Peers() + if len(peers) != 1 { + t.Fatalf("expected 1 peer, got %v", len(peers)) + } else if !peers[0].Synced() { + t.Fatal("fork peer that did not extend our chain should be marked synced") + } + + // verify s1 stayed on its own chain + if cm1.Tip() != s1Tip { + t.Fatalf("s1 tip should not have changed: expected %v, got %v", s1Tip, cm1.Tip()) + } + + // our tip moving is not a reason to ask the peer again; only a relay from + // them (which triggers a resync) is + testutil.MineBlocks(t, cm1, types.VoidAddress, 1) + time.Sleep(time.Second) // ~10 sync intervals + if reqs := fm.headerReqs.Load(); reqs != settledHeaders { + t.Fatalf("s1 re-synced the fork peer after its own tip moved: %v more header requests", reqs-settledHeaders) + } else if reqs := fm.blockReqs.Load(); reqs != settled { + t.Fatalf("s1 re-downloaded the fork: %v extra block requests", reqs-settled) + } +} + +// TestSyncAcrossBatches verifies that a peer whose blocks do extend our chain +// keeps being synced from, even when it has more headers than a single batch +// can carry. It guards the "nothing for us" check in syncLoop against marking a +// useful peer as synced. +func TestSyncAcrossBatches(t *testing.T) { + log := zaptest.NewLogger(t) + + // s1 starts at genesis + s1, cm1 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer1"))) + defer s1.Close() + + // s2 has 20 blocks but only serves 5 headers at a time, so it takes + // several batches (and always reports headers remaining) to catch up + n, genesis := testutil.Network() + store2, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) + if err != nil { + t.Fatal(err) + } + fm := &forkManager{headerLimit: 5, Manager: chain.NewManager(store2)} + testutil.MineBlocks(t, fm.Manager, types.VoidAddress, 20) + + l2, err := net.Listen("tcp", ":0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { l2.Close() }) + + s2 := syncer.New(l2, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ + GenesisID: genesis.ID(), + UniqueID: gateway.GenerateUniqueID(), + NetAddress: l2.Addr().String(), + }, syncer.WithSyncInterval(time.Hour)) // effectively disabled + go s2.Run() + defer s2.Close() + + if _, err := s1.Connect(context.Background(), s2.Addr()); err != nil { + t.Fatal(err) + } + + want := fm.Manager.Tip() + for range 100 { + if cm1.Tip() == want { + return + } + time.Sleep(100 * time.Millisecond) + } + t.Fatalf("s1 did not catch up: expected %v, got %v", want, cm1.Tip()) +} From 6601978340b04fd885883c66cf1adee561f5e4fc Mon Sep 17 00:00:00 2001 From: Chris Schinnerl <3903476+ChrisSchinnerl@users.noreply.github.com> Date: Tue, 22 Sep 2026 13:30:21 +0200 Subject: [PATCH 2/5] add TestHeavierForkAcrossBatches --- syncer/syncer_test.go | 51 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 51 insertions(+) diff --git a/syncer/syncer_test.go b/syncer/syncer_test.go index 72d0bf65..7978a472 100644 --- a/syncer/syncer_test.go +++ b/syncer/syncer_test.go @@ -1021,3 +1021,54 @@ func TestSyncAcrossBatches(t *testing.T) { } t.Fatalf("s1 did not catch up: expected %v, got %v", want, cm1.Tip()) } + +// TestHeavierForkAcrossBatches verifies that we adopt a peer's chain when the +// chain as a whole is heavier than ours, but its first batch of headers is not. +// A batch that leaves our tip unchanged does not mean the peer has nothing for +// us: it still reports headers remaining, and its continuation outweighs our +// chain. +func TestHeavierForkAcrossBatches(t *testing.T) { + log := zaptest.NewLogger(t) + + // s1 has a short chain of its own + s1, cm1 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer1"))) + defer s1.Close() + testutil.MineBlocks(t, cm1, types.VoidAddress, 10) + + // s2 is on a heavier fork, but only serves 5 headers at a time, so its + // first batch is shorter than s1's chain and does not trigger a reorg + n, genesis := testutil.Network() + store2, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) + if err != nil { + t.Fatal(err) + } + fm := &forkManager{headerLimit: 5, Manager: chain.NewManager(store2)} + testutil.MineBlocks(t, fm.Manager, types.Address{1}, 30) + + l2, err := net.Listen("tcp", ":0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { l2.Close() }) + + s2 := syncer.New(l2, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ + GenesisID: genesis.ID(), + UniqueID: gateway.GenerateUniqueID(), + NetAddress: l2.Addr().String(), + }, syncer.WithSyncInterval(time.Hour)) // effectively disabled + go s2.Run() + defer s2.Close() + + if _, err := s1.Connect(context.Background(), s2.Addr()); err != nil { + t.Fatal(err) + } + + want := fm.Manager.Tip() + for range 100 { + if cm1.Tip() == want { + return + } + time.Sleep(100 * time.Millisecond) + } + t.Fatalf("s1 did not adopt the heavier chain: expected %v, got %v", want, cm1.Tip()) +} From 31c8cc482b70a410b386b4f990e3017342e66709 Mon Sep 17 00:00:00 2001 From: Chris Schinnerl <3903476+ChrisSchinnerl@users.noreply.github.com> Date: Tue, 22 Sep 2026 13:42:05 +0200 Subject: [PATCH 3/5] use chain weight instead --- ...k_before_downloading_blocks_from_a_peer.md | 7 + ...ose_blocks_do_not_extend_our_best_chain.md | 7 - syncer/peer.go | 10 +- syncer/syncer.go | 145 +++++++++----- syncer/syncer_test.go | 187 +++++++----------- 5 files changed, 181 insertions(+), 175 deletions(-) create mode 100644 .changeset/compared_header_chain_work_before_downloading_blocks_from_a_peer.md delete mode 100644 .changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md diff --git a/.changeset/compared_header_chain_work_before_downloading_blocks_from_a_peer.md b/.changeset/compared_header_chain_work_before_downloading_blocks_from_a_peer.md new file mode 100644 index 00000000..3bdf4eac --- /dev/null +++ b/.changeset/compared_header_chain_work_before_downloading_blocks_from_a_peer.md @@ -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. diff --git a/.changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md b/.changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md deleted file mode 100644 index 137e6d94..00000000 --- a/.changeset/stopped_repeatedly_syncing_from_peers_whose_blocks_do_not_extend_our_best_chain.md +++ /dev/null @@ -1,7 +0,0 @@ ---- -default: patch ---- - -# Stopped repeatedly syncing from peers whose blocks do not extend our best chain - -The syncer now skips a peer if our tip is unchanged since the last batch it served. Previously a peer stuck on a fork that would never outweigh our chain was queried every sync interval, re-downloading its entire fork indefinitely. diff --git a/syncer/peer.go b/syncer/peer.go index e763aaf0..3c1fbb1c 100644 --- a/syncer/peer.go +++ b/syncer/peer.go @@ -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{}) } } - 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. diff --git a/syncer/syncer.go b/syncer/syncer.go index ee445a82..c59649c1 100644 --- a/syncer/syncer.go +++ b/syncer/syncer.go @@ -781,6 +781,85 @@ func (s *Syncer) peerLoop(ctx context.Context) error { return nil } +// 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 + batches int // SendHeaders calls needed to reach tip + remaining uint64 // headers the peer has beyond tip + heavier bool +} + +// errPeerReorg is returned when a peer's chain changes mid-walk, leaving the +// index we were walking from off their best chain. The walk is abandoned and +// retried on the next tick rather than faulting the peer. +var errPeerReorg = errors.New("peer reorged during header walk") + +// 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. +// Only the running state is kept, so a peer on a chain we will never adopt +// costs a header walk and nothing else. +func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.State) (peerChain, 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 + } + fork, ok := s.cm.State(id) + if !ok { + return peerChain{}, errors.New("missing state for history") + } + headers, tip, remaining, err := p.SendHeaders(fork, 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 peerChain{}, err + } + pc := peerChain{fork: fork} + for { + if len(headers) > 0 { + pc.tip, pc.batches = headers[len(headers)-1], pc.batches+1 + } + pc.remaining = remaining + if tip.SufficientlyHeavierThan(ts) { + pc.heavier = true + return pc, nil + } else if remaining == 0 || len(headers) == 0 { + return pc, nil + } + if headers, tip, remaining, err = p.SendHeaders(tip, s.config.MaxSendHeaders, s.config.SendHeadersTimeout); err != nil { + if strings.Contains(err.Error(), "EOF") { + return peerChain{}, errPeerReorg + } + return peerChain{}, err + } + } + } + return peerChain{}, errors.New("no common history") +} + +// syncPeerChain re-requests pc's headers from p and downloads each batch's +// blocks as it arrives, so only one batch is held in memory at a time. It +// returns what the peer actually served, which may fall short of pc if the peer +// reorged since the walk. +func (s *Syncer) syncPeerChain(ctx context.Context, p *Peer, pc peerChain) (peerChain, error) { + cs, synced := pc.fork, peerChain{fork: pc.fork} + for range pc.batches { + headers, tip, remaining, err := p.SendHeaders(cs, s.config.MaxSendHeaders, s.config.SendHeadersTimeout) + if err != nil { + return synced, err + } else if len(headers) == 0 { + break // peer changed chains between the walk and now + } else if err := s.parallelSync(ctx, cs, headers); err != nil { + return synced, err + } + cs = tip + synced.tip, synced.batches, synced.remaining = headers[len(headers)-1], synced.batches+1, remaining + } + return synced, nil +} + func (s *Syncer) syncLoop(ctx context.Context) error { ticker := time.NewTicker(s.config.SyncInterval) defer ticker.Stop() @@ -801,73 +880,45 @@ 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, errPeerReorg) { + s.log.Debug("peer reorged during header walk", zap.Stringer("peer", r.peer)) + } else if r.err != nil { r.peer.setErr(r.err) - } else if len(r.headers) == 0 { + } else if !r.chain.heavier { + // their whole chain is lighter than ours; if they later learn + // of a block we don't have, relaying it 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))) - tip := s.cm.Tip() - 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)) + if synced, err := s.syncPeerChain(ctx, r.peer, r.chain); err != nil { s.log.Debug("sync failed", zap.Stringer("peer", r.peer), zap.Error(err)) - } else if s.cm.Tip() == tip { - // their blocks did not extend our best chain, so the peer - // has nothing for us. Mark them synced; if they later learn - // of a block we don't have, relaying it triggers a resync. - // - // NOTE: a peer stuck on a chain we will never adopt never - // reports remaining == 0, so without this we would ask them - // for the same headers every sync interval, indefinitely. - s.log.Debug("peer has no blocks that extend our chain", zap.Stringer("peer", r.peer), zap.Stringer("tip", tip)) - r.peer.setSynced(true) - } 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; relay their tip r.peer.setSynced(true) - go s.relayV2Header(r.headers[len(r.headers)-1], r.peer) + go s.relayV2Header(synced.tip, r.peer) } } } diff --git a/syncer/syncer_test.go b/syncer/syncer_test.go index 7978a472..8b756f55 100644 --- a/syncer/syncer_test.go +++ b/syncer/syncer_test.go @@ -326,7 +326,7 @@ func TestSendHeaders(t *testing.T) { if err != nil { t.Fatal(err) } - headers, rem, err := p.SendHeaders(cs, 90, time.Second) + headers, _, rem, err := p.SendHeaders(cs, 90, time.Second) if err != nil { t.Fatal(err) } else if len(headers) != 90 { @@ -863,8 +863,8 @@ func TestMaxInflightRPCsBackpressureNotDropped(t *testing.T) { } } -// forkManager serves a limited number of headers per request and counts block -// requests, simulating a peer on a fork longer than a single sync batch. +// forkManager serves a limited number of headers per request and counts +// requests, simulating a peer whose chain spans several sync batches. type forkManager struct { headerLimit uint64 headerReqs atomic.Uint64 @@ -882,44 +882,62 @@ func (fm *forkManager) BlocksForHistory(history []types.BlockID, maxBlocks uint6 return fm.Manager.BlocksForHistory(history, maxBlocks) } -// TestForkPeerNotResynced verifies that we stop syncing from a peer on a fork -// that never becomes our best chain. Such a peer is never marked synced, -// because it always has headers remaining; since our tip doesn't change, it -// would answer every subsequent sync round with the same headers, and we would -// re-download the same blocks indefinitely. -func TestForkPeerNotResynced(t *testing.T) { - log := zaptest.NewLogger(t) - - // s1 has the longer chain - s1, cm1 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer1"))) - defer s1.Close() - testutil.MineBlocks(t, cm1, types.VoidAddress, 20) +// newForkPeer starts a syncer serving at most headerLimit headers per request +// from a chain of the given length, mined to addr. Its own sync loop is +// disabled so it never adopts our chain. +func newForkPeer(t testing.TB, headerLimit uint64, addr types.Address, blocks int) (*syncer.Syncer, *forkManager) { + t.Helper() - // s2 is on a shorter fork, and only serves 5 headers at a time, so it - // always reports headers remaining n, genesis := testutil.Network() - store2, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) + store, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) if err != nil { t.Fatal(err) } - fm := &forkManager{headerLimit: 5, Manager: chain.NewManager(store2)} - // mine to a different address so the fork diverges at the genesis child, - // leaving more headers than a single batch can carry - testutil.MineBlocks(t, fm.Manager, types.Address{1}, 10) + fm := &forkManager{headerLimit: headerLimit, Manager: chain.NewManager(store)} + testutil.MineBlocks(t, fm.Manager, addr, blocks) - l2, err := net.Listen("tcp", ":0") + l, err := net.Listen("tcp", ":0") if err != nil { t.Fatal(err) } - t.Cleanup(func() { l2.Close() }) + t.Cleanup(func() { l.Close() }) - s2 := syncer.New(l2, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ + s := syncer.New(l, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ GenesisID: genesis.ID(), UniqueID: gateway.GenerateUniqueID(), - NetAddress: l2.Addr().String(), + NetAddress: l.Addr().String(), }, syncer.WithSyncInterval(time.Hour)) // effectively disabled - go s2.Run() - defer s2.Close() + go s.Run() + t.Cleanup(func() { s.Close() }) + return s, fm +} + +// waitForTip blocks until cm reaches want. +func waitForTip(t testing.TB, cm *chain.Manager, want types.ChainIndex, msg string) { + t.Helper() + for range 100 { + if cm.Tip() == want { + return + } + time.Sleep(100 * time.Millisecond) + } + t.Fatalf("%s: expected tip %v, got %v", msg, want, cm.Tip()) +} + +// TestForkPeerNotResynced verifies that we stop syncing from a peer on a fork +// that never becomes our best chain: we walk their headers once, download no +// blocks, and don't walk them again every sync interval. +func TestForkPeerNotResynced(t *testing.T) { + log := zaptest.NewLogger(t) + + // s1 has the longer chain + s1, cm1 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer1"))) + defer s1.Close() + testutil.MineBlocks(t, cm1, types.VoidAddress, 20) + + // s2 is on a shorter fork, serving 5 headers at a time so it always reports + // headers remaining + s2, fm := newForkPeer(t, 5, types.Address{1}, 10) s1Tip := cm1.Tip() @@ -927,29 +945,29 @@ func TestForkPeerNotResynced(t *testing.T) { t.Fatal(err) } - // wait for s1 to fetch a batch of fork blocks + // wait for s1 to walk the fork peer's headers for range 100 { - if fm.blockReqs.Load() > 0 { + if fm.headerReqs.Load() > 0 { break } time.Sleep(100 * time.Millisecond) } - if fm.blockReqs.Load() == 0 { - t.Fatal("s1 never synced any blocks from the fork peer") + if fm.headerReqs.Load() == 0 { + t.Fatal("s1 never requested headers from the fork peer") } - // s1 should give up on the peer rather than starting a fresh sync round - // every sync interval; allow a few in-flight rounds to settle + // allow any in-flight rounds to settle time.Sleep(time.Second) - settled, settledHeaders := fm.blockReqs.Load(), fm.headerReqs.Load() + settledHeaders := fm.headerReqs.Load() time.Sleep(2 * time.Second) // ~20 sync intervals if reqs := fm.headerReqs.Load(); reqs != settledHeaders { t.Fatalf("s1 kept starting sync rounds with the fork peer: %v more header requests", reqs-settledHeaders) - } else if reqs := fm.blockReqs.Load(); reqs != settled { - t.Fatalf("s1 kept re-syncing the fork peer: %v block requests, then %v more", settled, reqs-settled) } - // the peer had nothing that extends our chain, so it is done with + if reqs := fm.blockReqs.Load(); reqs != 0 { + t.Fatalf("s1 downloaded %v block batches from a chain that never outweighed ours", reqs) + } + peers := s1.Peers() if len(peers) != 1 { t.Fatalf("expected 1 peer, got %v", len(peers)) @@ -957,26 +975,22 @@ func TestForkPeerNotResynced(t *testing.T) { t.Fatal("fork peer that did not extend our chain should be marked synced") } - // verify s1 stayed on its own chain if cm1.Tip() != s1Tip { t.Fatalf("s1 tip should not have changed: expected %v, got %v", s1Tip, cm1.Tip()) } - // our tip moving is not a reason to ask the peer again; only a relay from - // them (which triggers a resync) is + // our tip moving is not a reason to ask again; only a relay from them is testutil.MineBlocks(t, cm1, types.VoidAddress, 1) time.Sleep(time.Second) // ~10 sync intervals if reqs := fm.headerReqs.Load(); reqs != settledHeaders { t.Fatalf("s1 re-synced the fork peer after its own tip moved: %v more header requests", reqs-settledHeaders) - } else if reqs := fm.blockReqs.Load(); reqs != settled { - t.Fatalf("s1 re-downloaded the fork: %v extra block requests", reqs-settled) + } else if reqs := fm.blockReqs.Load(); reqs != 0 { + t.Fatalf("s1 downloaded %v block batches after its own tip moved", reqs) } } -// TestSyncAcrossBatches verifies that a peer whose blocks do extend our chain -// keeps being synced from, even when it has more headers than a single batch -// can carry. It guards the "nothing for us" check in syncLoop against marking a -// useful peer as synced. +// TestSyncAcrossBatches verifies that we keep syncing from a peer with more +// headers than a single batch can carry. func TestSyncAcrossBatches(t *testing.T) { log := zaptest.NewLogger(t) @@ -984,49 +998,18 @@ func TestSyncAcrossBatches(t *testing.T) { s1, cm1 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer1"))) defer s1.Close() - // s2 has 20 blocks but only serves 5 headers at a time, so it takes - // several batches (and always reports headers remaining) to catch up - n, genesis := testutil.Network() - store2, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) - if err != nil { - t.Fatal(err) - } - fm := &forkManager{headerLimit: 5, Manager: chain.NewManager(store2)} - testutil.MineBlocks(t, fm.Manager, types.VoidAddress, 20) - - l2, err := net.Listen("tcp", ":0") - if err != nil { - t.Fatal(err) - } - t.Cleanup(func() { l2.Close() }) - - s2 := syncer.New(l2, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ - GenesisID: genesis.ID(), - UniqueID: gateway.GenerateUniqueID(), - NetAddress: l2.Addr().String(), - }, syncer.WithSyncInterval(time.Hour)) // effectively disabled - go s2.Run() - defer s2.Close() + // s2 has 20 blocks but serves 5 headers at a time, so catching up takes + // several batches + s2, fm := newForkPeer(t, 5, types.VoidAddress, 20) if _, err := s1.Connect(context.Background(), s2.Addr()); err != nil { t.Fatal(err) } - - want := fm.Manager.Tip() - for range 100 { - if cm1.Tip() == want { - return - } - time.Sleep(100 * time.Millisecond) - } - t.Fatalf("s1 did not catch up: expected %v, got %v", want, cm1.Tip()) + waitForTip(t, cm1, fm.Manager.Tip(), "s1 did not catch up") } -// TestHeavierForkAcrossBatches verifies that we adopt a peer's chain when the -// chain as a whole is heavier than ours, but its first batch of headers is not. -// A batch that leaves our tip unchanged does not mean the peer has nothing for -// us: it still reports headers remaining, and its continuation outweighs our -// chain. +// TestHeavierForkAcrossBatches verifies that we adopt a peer's chain that is +// heavier than ours overall, even though its first batch of headers is not. func TestHeavierForkAcrossBatches(t *testing.T) { log := zaptest.NewLogger(t) @@ -1035,40 +1018,12 @@ func TestHeavierForkAcrossBatches(t *testing.T) { defer s1.Close() testutil.MineBlocks(t, cm1, types.VoidAddress, 10) - // s2 is on a heavier fork, but only serves 5 headers at a time, so its - // first batch is shorter than s1's chain and does not trigger a reorg - n, genesis := testutil.Network() - store2, err := chain.NewDBStore(chain.NewMemDB(), n, genesis, nil) - if err != nil { - t.Fatal(err) - } - fm := &forkManager{headerLimit: 5, Manager: chain.NewManager(store2)} - testutil.MineBlocks(t, fm.Manager, types.Address{1}, 30) - - l2, err := net.Listen("tcp", ":0") - if err != nil { - t.Fatal(err) - } - t.Cleanup(func() { l2.Close() }) - - s2 := syncer.New(l2, fm, testutil.NewEphemeralPeerStore(), gateway.Header{ - GenesisID: genesis.ID(), - UniqueID: gateway.GenerateUniqueID(), - NetAddress: l2.Addr().String(), - }, syncer.WithSyncInterval(time.Hour)) // effectively disabled - go s2.Run() - defer s2.Close() + // s2's fork is heavier, but its first 5-header batch is shorter than s1's + // chain and does not trigger a reorg + s2, fm := newForkPeer(t, 5, types.Address{1}, 30) if _, err := s1.Connect(context.Background(), s2.Addr()); err != nil { t.Fatal(err) } - - want := fm.Manager.Tip() - for range 100 { - if cm1.Tip() == want { - return - } - time.Sleep(100 * time.Millisecond) - } - t.Fatalf("s1 did not adopt the heavier chain: expected %v, got %v", want, cm1.Tip()) + waitForTip(t, cm1, fm.Manager.Tip(), "s1 did not adopt the heavier chain") } From 523fc29a731c43559fab234ede70a0541e131740 Mon Sep 17 00:00:00 2001 From: Chris Schinnerl <3903476+ChrisSchinnerl@users.noreply.github.com> Date: Wed, 23 Sep 2026 11:30:14 +0200 Subject: [PATCH 4/5] address comments --- syncer/syncer.go | 128 +++++++++++++++++++++++++++++++----------- syncer/syncer_test.go | 18 ++++-- 2 files changed, 109 insertions(+), 37 deletions(-) diff --git a/syncer/syncer.go b/syncer/syncer.go index c59649c1..b5bec14b 100644 --- a/syncer/syncer.go +++ b/syncer/syncer.go @@ -781,25 +781,65 @@ 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 - batches int // SendHeaders calls needed to reach tip + 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 } -// errPeerReorg is returned when a peer's chain changes mid-walk, leaving the -// index we were walking from off their best chain. The walk is abandoned and -// retried on the next tick rather than faulting the peer. -var errPeerReorg = errors.New("peer reorged during header walk") +// 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 + maxHeaderWalk = 2 * 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) + errWalkTimeout = fmt.Errorf("%w: exceeded %v", errWalkAbandoned, maxHeaderWalk) + + // 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") +) // 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. -// Only the running state is kept, so a peer on a chain we will never adopt -// costs a header walk and nothing else. func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.State) (peerChain, error) { + // 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) { + // don't fault the peer for our own 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 @@ -810,16 +850,22 @@ func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.Sta if !ok { return peerChain{}, errors.New("missing state for history") } - headers, tip, remaining, err := p.SendHeaders(fork, s.config.MaxSendHeaders, s.config.SendHeadersTimeout) + 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 := peerChain{fork: fork} + pc, cs := peerChain{fork: fork}, fork for { if len(headers) > 0 { - pc.tip, pc.batches = headers[len(headers)-1], pc.batches+1 + 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 + } } pc.remaining = remaining if tip.SufficientlyHeavierThan(ts) { @@ -828,7 +874,8 @@ func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.Sta } else if remaining == 0 || len(headers) == 0 { return pc, nil } - if headers, tip, remaining, err = p.SendHeaders(tip, s.config.MaxSendHeaders, s.config.SendHeadersTimeout); err != nil { + cs = tip + if headers, tip, remaining, err = sendHeaders(cs); err != nil { if strings.Contains(err.Error(), "EOF") { return peerChain{}, errPeerReorg } @@ -839,23 +886,38 @@ func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.Sta return peerChain{}, errors.New("no common history") } -// syncPeerChain re-requests pc's headers from p and downloads each batch's -// blocks as it arrives, so only one batch is held in memory at a time. It -// returns what the peer actually served, which may fall short of pc if the peer -// reorged since the walk. -func (s *Syncer) syncPeerChain(ctx context.Context, p *Peer, pc peerChain) (peerChain, error) { - cs, synced := pc.fork, peerChain{fork: pc.fork} - for range pc.batches { +// 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) { + 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 { - return synced, err - } else if len(headers) == 0 { - break // peer changed chains between the walk and now + 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 { return synced, err } cs = tip - synced.tip, synced.batches, synced.remaining = headers[len(headers)-1], synced.batches+1, remaining + synced.tip, synced.remaining = headers[len(headers)-1], remaining + synced.batchTips = append(synced.batchTips, synced.tip.ID()) } return synced, nil } @@ -899,24 +961,26 @@ func (s *Syncer) syncLoop(ctx context.Context) error { // sync each set of headers as they arrive seen := make(map[types.BlockID]bool) for range peers { - if r := <-respChan; errors.Is(r.err, errPeerReorg) { - s.log.Debug("peer reorged during header walk", zap.Stringer("peer", r.peer)) + if r := <-respChan; errors.Is(r.err, errWalkAbandoned) { + 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 !r.chain.heavier { - // their whole chain is lighter than ours; if they later learn - // of a block we don't have, relaying it triggers a resync - s.log.Debug("peer chain does not outweigh ours", zap.Stringer("peer", r.peer), zap.Int("batches", r.chain.batches)) + // 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.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.chain.fork.Index), zap.Int("batches", r.chain.batches)) - if synced, err := s.syncPeerChain(ctx, r.peer, r.chain); 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 synced.batches > 0 && synced.remaining == 0 { - // peer sent all their headers; relay their tip + } else if synced.batches() > 0 && synced.remaining == 0 { + // peer sent all their headers r.peer.setSynced(true) go s.relayV2Header(synced.tip, r.peer) } diff --git a/syncer/syncer_test.go b/syncer/syncer_test.go index 8b756f55..499b9e61 100644 --- a/syncer/syncer_test.go +++ b/syncer/syncer_test.go @@ -107,8 +107,9 @@ func TestSyncer(t *testing.T) { s2, cm2 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer2"))) defer s2.Close() - // mine enough blocks to test both v1 and v2 regimes - testutil.MineBlocks(t, cm1, types.VoidAddress, int(cm1.TipState().Network.HardforkV2.RequireHeight+100)) + // mine enough blocks to cover both v1 and v2 regimes, and to split the + // download into enough requests that the bad peer is not starved of work + testutil.MineBlocks(t, cm1, types.VoidAddress, int(cm1.TipState().Network.HardforkV2.RequireHeight+1000)) if _, err := s1.Connect(context.Background(), s2.Addr()); err != nil { t.Fatal(err) @@ -151,8 +152,9 @@ func TestSyncWithBadPeer(t *testing.T) { s2, cm2 := newTestSyncer(t, syncer.WithLogger(log.Named("syncer2"))) defer s2.Close() - // mine enough blocks to test both v1 and v2 regimes - testutil.MineBlocks(t, cm1, types.VoidAddress, int(cm1.TipState().Network.HardforkV2.RequireHeight+100)) + // mine enough blocks to cover both v1 and v2 regimes, and to split the + // download into enough requests that the bad peer is not starved of work + testutil.MineBlocks(t, cm1, types.VoidAddress, int(cm1.TipState().Network.HardforkV2.RequireHeight+1000)) // simulate another peer, one that returns invalid blocks _, genesis := testutil.Network() @@ -186,7 +188,13 @@ func TestSyncWithBadPeer(t *testing.T) { if cm1.Tip() != cm2.Tip() { t.Fatalf("tips are not equal: %v != %v", cm1.Tip(), cm2.Tip()) } - // bad peer should be banned + // the ban propagates asynchronously, so wait for the peer to be dropped + for range 100 { + if len(s2.Peers()) == 1 { + break + } + time.Sleep(100 * time.Millisecond) + } if peers := s2.Peers(); len(peers) != 1 { t.Fatalf("expected 1 peer, got %v", peers) } else if peers[0].UniqueID() == badID { From 6679b8024cd06fa6258dfa9427775d3354332251 Mon Sep 17 00:00:00 2001 From: Chris Schinnerl <3903476+ChrisSchinnerl@users.noreply.github.com> Date: Fri, 25 Sep 2026 11:19:59 +0200 Subject: [PATCH 5/5] add limit to syncChain and attribute timeout errors to peer --- syncer/syncer.go | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/syncer/syncer.go b/syncer/syncer.go index b5bec14b..0c907d31 100644 --- a/syncer/syncer.go +++ b/syncer/syncer.go @@ -807,8 +807,13 @@ const ( // retainedBatches bounds how many batches a walk keeps retainedBatches = 2 - // maxHeaderWalk bounds one peer's walk; a sync round waits for every peer + // 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 ( @@ -816,11 +821,11 @@ var ( // peer stays eligible and is retried next tick errWalkAbandoned = errors.New("header walk abandoned") errPeerReorg = fmt.Errorf("%w: peer reorged", errWalkAbandoned) - errWalkTimeout = fmt.Errorf("%w: exceeded %v", errWalkAbandoned, maxHeaderWalk) // 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 @@ -835,7 +840,6 @@ func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.Sta } headers, tip, remaining, err := p.SendHeaders(cs, s.config.MaxSendHeaders, timeout) if err != nil && !time.Now().Before(deadline) { - // don't fault the peer for our own deadline return nil, consensus.State{}, 0, errWalkTimeout } return headers, tip, remaining, err @@ -890,6 +894,9 @@ func (s *Syncer) walkPeerChain(p *Peer, hist [32]types.BlockID, ts consensus.Sta // 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 {