Skip to content

Commit aae7e66

Browse files
authored
eth: retry sync after peer-set replacements (#2391)
1 parent 321a24a commit aae7e66

4 files changed

Lines changed: 265 additions & 5 deletions

File tree

eth/peerset.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,7 @@ type peerSet struct {
6060
peers map[string]*ethPeer // Peers connected on the `eth` protocol
6161
snapPeers int // Number of `snap` compatible peers for connection prioritization
6262
witPeers int // Number of `wit` compatible peers for connection prioritization
63+
revision uint64 // Peer-set revision
6364

6465
snapWait map[string]chan *snap.Peer // Peers connected on `eth` waiting for their snap extension
6566
snapPend map[string]*snap.Peer // Peers connected on the `snap` protocol, but not yet on `eth`
@@ -267,6 +268,7 @@ func (ps *peerSet) registerPeer(peer *eth.Peer, extSnap *snap.Peer, extWit *wit.
267268
}
268269

269270
ps.peers[id] = eth
271+
ps.revision++
270272

271273
return nil
272274
}
@@ -283,6 +285,7 @@ func (ps *peerSet) unregisterPeer(id string) error {
283285
}
284286

285287
delete(ps.peers, id)
288+
ps.revision++
286289

287290
if peer.snapExt != nil {
288291
ps.snapPeers--
@@ -448,6 +451,13 @@ func (ps *peerSet) len() int {
448451
return len(ps.peers)
449452
}
450453

454+
func (ps *peerSet) currentRevision() uint64 {
455+
ps.lock.RLock()
456+
defer ps.lock.RUnlock()
457+
458+
return ps.revision
459+
}
460+
451461
// snapLen returns if the current number of `snap` peers in the set.
452462
func (ps *peerSet) snapLen() int {
453463
ps.lock.RLock()

eth/peerset_test.go

Lines changed: 206 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,18 +2,224 @@ package eth
22

33
import (
44
"crypto/rand"
5+
"errors"
56
"math/big"
67
"testing"
78
"time"
89

910
"github.com/ethereum/go-ethereum/common"
1011
"github.com/ethereum/go-ethereum/eth/protocols/eth"
12+
"github.com/ethereum/go-ethereum/eth/protocols/snap"
1113
"github.com/ethereum/go-ethereum/eth/protocols/wit"
1214
"github.com/ethereum/go-ethereum/log"
1315
"github.com/ethereum/go-ethereum/p2p"
1416
"github.com/ethereum/go-ethereum/p2p/enode"
1517
)
1618

19+
func TestPeerSetExtensions(t *testing.T) {
20+
t.Run("snap", func(t *testing.T) {
21+
caps := []p2p.Cap{{Name: eth.ProtocolName, Version: eth.ETH69}, {Name: snap.ProtocolName, Version: snap.SNAP1}}
22+
testPeerSetExtension(t, caps, false, errSnapWithoutEth, newPeerSetSnapPeer, (*peerSet).registerSnapExtension, (*peerSet).waitSnapExtension)
23+
})
24+
t.Run("wit", func(t *testing.T) {
25+
caps := []p2p.Cap{{Name: eth.ProtocolName, Version: eth.ETH69}, {Name: wit.ProtocolName, Version: wit.WIT2}}
26+
testPeerSetExtension(t, caps, true, errWitWithoutEth, newPeerSetWitPeer, (*peerSet).registerWitExtension, (*peerSet).waitWitExtension)
27+
})
28+
}
29+
30+
func TestPeerSetRevision(t *testing.T) {
31+
ps := newPeerSet()
32+
peer := newPeerSetEthPeer(t, enode.ID{1}, nil)
33+
34+
if err := ps.unregisterPeer(peer.ID()); !errors.Is(err, errPeerNotRegistered) {
35+
t.Fatalf("unexpected unregister error: %v", err)
36+
}
37+
if revision := ps.currentRevision(); revision != 0 {
38+
t.Fatalf("failed unregister changed revision to %d", revision)
39+
}
40+
if err := ps.registerPeer(peer, nil, nil); err != nil {
41+
t.Fatal(err)
42+
}
43+
if err := ps.registerPeer(peer, nil, nil); !errors.Is(err, errPeerAlreadyRegistered) {
44+
t.Fatalf("unexpected duplicate registration error: %v", err)
45+
}
46+
if revision := ps.currentRevision(); revision != 1 {
47+
t.Fatalf("registration changed revision to %d", revision)
48+
}
49+
if err := ps.unregisterPeer(peer.ID()); err != nil {
50+
t.Fatal(err)
51+
}
52+
if revision := ps.currentRevision(); revision != 2 {
53+
t.Fatalf("unregister changed revision to %d", revision)
54+
}
55+
ps.close()
56+
if err := ps.registerPeer(peer, nil, nil); !errors.Is(err, errPeerSetClosed) {
57+
t.Fatalf("unexpected closed set error: %v", err)
58+
}
59+
if revision := ps.currentRevision(); revision != 2 {
60+
t.Fatalf("failed registration changed revision to %d", revision)
61+
}
62+
}
63+
64+
type extensionResult[T comparable] struct {
65+
peer T
66+
err error
67+
}
68+
69+
func testPeerSetExtension[T comparable](
70+
t *testing.T,
71+
caps []p2p.Cap,
72+
witness bool,
73+
incompatibleError error,
74+
newExtension func(*testing.T, enode.ID, []p2p.Cap) T,
75+
register func(*peerSet, T) error,
76+
wait func(*peerSet, *eth.Peer) (T, error),
77+
) {
78+
t.Helper()
79+
80+
t.Run("extension first", func(t *testing.T) {
81+
ps := newPeerSet()
82+
defer ps.close()
83+
84+
id := enode.ID{1}
85+
main := newPeerSetEthPeer(t, id, caps)
86+
ext := newExtension(t, id, caps)
87+
if err := register(ps, ext); err != nil {
88+
t.Fatal(err)
89+
}
90+
if err := register(ps, ext); !errors.Is(err, errPeerAlreadyRegistered) {
91+
t.Fatalf("unexpected duplicate registration error: %v", err)
92+
}
93+
got, err := wait(ps, main)
94+
if err != nil || got != ext {
95+
t.Fatalf("extension mismatch: got %v, want %v, err %v", got, ext, err)
96+
}
97+
})
98+
99+
t.Run("main first", func(t *testing.T) {
100+
ps := newPeerSet()
101+
defer ps.close()
102+
103+
id := enode.ID{2}
104+
main := newPeerSetEthPeer(t, id, caps)
105+
ext := newExtension(t, id, caps)
106+
result := make(chan extensionResult[T], 1)
107+
go func() {
108+
peer, err := wait(ps, main)
109+
result <- extensionResult[T]{peer, err}
110+
}()
111+
waitForPeerSetWaiter(t, ps, id.String(), witness)
112+
if err := register(ps, ext); err != nil {
113+
t.Fatal(err)
114+
}
115+
if got := <-result; got.err != nil || got.peer != ext {
116+
t.Fatalf("extension mismatch: got %v, want %v, err %v", got.peer, ext, got.err)
117+
}
118+
})
119+
120+
t.Run("closed", func(t *testing.T) {
121+
ps := newPeerSet()
122+
id := enode.ID{3}
123+
main := newPeerSetEthPeer(t, id, caps)
124+
result := make(chan extensionResult[T], 1)
125+
go func() {
126+
peer, err := wait(ps, main)
127+
result <- extensionResult[T]{peer, err}
128+
}()
129+
waitForPeerSetWaiter(t, ps, id.String(), witness)
130+
ps.close()
131+
if got := <-result; !errors.Is(got.err, errPeerSetClosed) {
132+
t.Fatalf("unexpected close error: %v", got.err)
133+
}
134+
})
135+
136+
t.Run("registered", func(t *testing.T) {
137+
ps := newPeerSet()
138+
defer ps.close()
139+
140+
id := enode.ID{4}
141+
main := newPeerSetEthPeer(t, id, caps)
142+
ext := newExtension(t, id, caps)
143+
if err := ps.registerPeer(main, nil, nil); err != nil {
144+
t.Fatal(err)
145+
}
146+
if err := register(ps, ext); !errors.Is(err, errPeerAlreadyRegistered) {
147+
t.Fatalf("unexpected extension registration error: %v", err)
148+
}
149+
if _, err := wait(ps, main); !errors.Is(err, errPeerAlreadyRegistered) {
150+
t.Fatalf("unexpected extension wait error: %v", err)
151+
}
152+
})
153+
154+
t.Run("incompatible", func(t *testing.T) {
155+
ps := newPeerSet()
156+
defer ps.close()
157+
158+
id := enode.ID{5}
159+
ext := newExtension(t, id, caps[1:])
160+
if err := register(ps, ext); !errors.Is(err, incompatibleError) {
161+
t.Fatalf("unexpected extension registration error: %v", err)
162+
}
163+
main := newPeerSetEthPeer(t, id, caps[:1])
164+
var zero T
165+
if got, err := wait(ps, main); err != nil || got != zero {
166+
t.Fatalf("unexpected extension: got %v, err %v", got, err)
167+
}
168+
})
169+
}
170+
171+
func newPeerSetEthPeer(t *testing.T, id enode.ID, caps []p2p.Cap) *eth.Peer {
172+
t.Helper()
173+
peer, rw := newPeerSetProtocolPeer(t, id, caps)
174+
result := eth.NewPeer(eth.ETH69, peer, rw, nil)
175+
t.Cleanup(result.Close)
176+
return result
177+
}
178+
179+
func newPeerSetSnapPeer(t *testing.T, id enode.ID, caps []p2p.Cap) *snap.Peer {
180+
t.Helper()
181+
peer, rw := newPeerSetProtocolPeer(t, id, caps)
182+
return snap.NewPeer(snap.SNAP1, peer, rw)
183+
}
184+
185+
func newPeerSetWitPeer(t *testing.T, id enode.ID, caps []p2p.Cap) *wit.Peer {
186+
t.Helper()
187+
peer, rw := newPeerSetProtocolPeer(t, id, caps)
188+
result := wit.NewPeer(wit.WIT2, peer, rw, log.New())
189+
t.Cleanup(result.Close)
190+
return result
191+
}
192+
193+
func newPeerSetProtocolPeer(t *testing.T, id enode.ID, caps []p2p.Cap) (*p2p.Peer, p2p.MsgReadWriter) {
194+
t.Helper()
195+
app, net := p2p.MsgPipe()
196+
t.Cleanup(func() {
197+
app.Close()
198+
net.Close()
199+
})
200+
return p2p.NewPeer(id, "test", caps), net
201+
}
202+
203+
func waitForPeerSetWaiter(t *testing.T, ps *peerSet, id string, witness bool) {
204+
t.Helper()
205+
deadline := time.Now().Add(time.Second)
206+
for time.Now().Before(deadline) {
207+
ps.lock.RLock()
208+
_, snapReady := ps.snapWait[id]
209+
_, witReady := ps.witWait[id]
210+
ps.lock.RUnlock()
211+
ready := snapReady
212+
if witness {
213+
ready = witReady
214+
}
215+
if ready {
216+
return
217+
}
218+
time.Sleep(time.Millisecond)
219+
}
220+
t.Fatal("extension waiter was not registered")
221+
}
222+
17223
func TestPeerSetForgetTransactions(t *testing.T) {
18224
t.Parallel()
19225

eth/sync.go

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -61,8 +61,8 @@ type chainSyncer struct {
6161
peerEventCh chan struct{}
6262
doneCh chan error // non-nil when sync is running
6363

64-
peersUnavailableUntil time.Time
65-
peersUnavailableAtCount int
64+
peersUnavailableUntil time.Time
65+
observedPeerRevision uint64
6666
}
6767

6868
// chainSyncOp is a scheduled sync operation.
@@ -112,6 +112,7 @@ func (cs *chainSyncer) loop() {
112112
retry := newResettableTimer()
113113
defer retry.stop()
114114

115+
cs.observedPeerRevision = cs.handler.peers.currentRevision()
115116
for {
116117
if op, wait := cs.nextSyncOp(); op != nil {
117118
retry.stop()
@@ -143,7 +144,6 @@ func (cs *chainSyncer) onSyncDone(err error) {
143144

144145
if errors.Is(err, downloader.ErrPeersUnavailable) || errors.Is(err, downloader.ErrPeerBackedOff) || errors.Is(err, whitelist.ErrNoRemote) {
145146
cs.peersUnavailableUntil = time.Now().Add(forceSyncCycle)
146-
cs.peersUnavailableAtCount = cs.handler.peers.len()
147147
} else {
148148
cs.peersUnavailableUntil = time.Time{}
149149
}
@@ -159,9 +159,11 @@ func (cs *chainSyncer) onSyncDone(err error) {
159159
}
160160

161161
func (cs *chainSyncer) onPeerEvent() {
162-
if !cs.peersUnavailableUntil.IsZero() && cs.handler.peers.len() != cs.peersUnavailableAtCount {
162+
revision := cs.handler.peers.currentRevision()
163+
if !cs.peersUnavailableUntil.IsZero() && revision != cs.observedPeerRevision {
163164
cs.peersUnavailableUntil = time.Time{}
164165
}
166+
cs.observedPeerRevision = revision
165167
}
166168

167169
func (cs *chainSyncer) shutdown() {

eth/sync_test.go

Lines changed: 43 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -195,6 +195,7 @@ func TestChainSyncerCooldownSurvivesBlockAnnounce(t *testing.T) {
195195
if err := handler.downloader.RegisterPeer(peer.ID(), eth.ETH68, &ethPeer{Peer: peer}); err != nil {
196196
t.Fatal(err)
197197
}
198+
syncer.onPeerEvent()
198199

199200
syncer.onSyncDone(downloader.ErrPeersUnavailable)
200201
if syncer.peersUnavailableUntil.IsZero() {
@@ -210,9 +211,50 @@ func TestChainSyncerCooldownSurvivesBlockAnnounce(t *testing.T) {
210211
if err := handler.downloader.RegisterPeer(peer2.ID(), eth.ETH68, &ethPeer{Peer: peer2}); err != nil {
211212
t.Fatal(err)
212213
}
214+
if err := handler.peers.unregisterPeer(peer.ID()); err != nil {
215+
t.Fatal(err)
216+
}
217+
if handler.peers.len() != 1 {
218+
t.Fatalf("peer replacement should preserve peer count, have %d", handler.peers.len())
219+
}
220+
syncer.onSyncDone(downloader.ErrPeersUnavailable)
213221
syncer.onPeerEvent()
214222
if !syncer.peersUnavailableUntil.IsZero() {
215-
t.Fatal("a genuine peer-set change must clear the cooldown")
223+
t.Fatal("a peer replacement must clear the cooldown")
224+
}
225+
}
226+
227+
func TestChainSyncerLoopInitializesPeerRevision(t *testing.T) {
228+
handler, cleanup := newChainSyncerTestHandler(t)
229+
defer cleanup()
230+
handler.maxPeers = defaultMinSyncPeers
231+
232+
peer := registerPeerWithTD(t, handler.peers, 1_000_000)
233+
if err := handler.downloader.RegisterPeer(peer.ID(), eth.ETH68, &ethPeer{Peer: peer}); err != nil {
234+
t.Fatal(err)
235+
}
236+
237+
syncer := handler.chainSync
238+
syncer.peersUnavailableUntil = time.Now().Add(time.Hour)
239+
handler.wg.Add(1)
240+
done := make(chan struct{})
241+
go func() {
242+
syncer.loop()
243+
close(done)
244+
}()
245+
246+
if !syncer.handlePeerEvent() {
247+
t.Fatal("chain syncer stopped before processing the peer event")
248+
}
249+
close(handler.quitSync)
250+
251+
select {
252+
case <-done:
253+
case <-time.After(time.Second):
254+
t.Fatal("chain syncer did not stop")
255+
}
256+
if syncer.peersUnavailableUntil.IsZero() {
257+
t.Fatal("an initial peer event with no peer-set change must not clear the cooldown")
216258
}
217259
}
218260

0 commit comments

Comments
 (0)