@cryptotaxi247 / kubo / commits / 94de60557

startup bootstrap

Will Scott committed Apr 14, 2020 at 12:26 UTC 94de6055702618c9b1b5f790965e18d72fd75434
1 file changed +40 -11
test/integration/wan_lan_dht_test.go
+40 -11
@@ -7,6 +7,7 @@ import (
7 "math"
8 "math/rand"
9 "net"
10 + "os"
11 "testing"
12 "time"
13
@@ -18,6 +19,7 @@ import (
19 corenet "github.com/libp2p/go-libp2p-core/network"
20 peer "github.com/libp2p/go-libp2p-core/peer"
21 "github.com/libp2p/go-libp2p-core/peerstore"
22 + kbucket "github.com/libp2p/go-libp2p-kbucket"
23 testutil "github.com/libp2p/go-libp2p-testing/net"
24 mocknet "github.com/libp2p/go-libp2p/p2p/net/mock"
25
@@ -90,6 +92,8 @@ func RunDHTConnectivity(conf testutil.LatencyConfig, numPeers int) error {
92 wanPeers := []*core.IpfsNode{}
93 lanPeers := []*core.IpfsNode{}
94
95 + connectionContext, connCtxCancel := context.WithTimeout(ctx, 15*time.Second)
96 + defer connCtxCancel()
97 for i := 0; i < numPeers; i++ {
98 wanPeer, err := core.NewNode(ctx, &core.BuildCfg{
99 Online: true,
@@ -103,7 +107,7 @@ func RunDHTConnectivity(conf testutil.LatencyConfig, numPeers int) error {
107 wanPeer.Peerstore.AddAddr(wanPeer.Identity, wanAddr, peerstore.PermanentAddrTTL)
108 for _, p := range wanPeers {
109 mn.LinkPeers(p.Identity, wanPeer.Identity)
106 - mn.ConnectPeers(p.Identity, wanPeer.Identity)
110 + wanPeer.PeerHost.Connect(connectionContext, p.Peerstore.PeerInfo(p.Identity))
111 }
112 wanPeers = append(wanPeers, wanPeer)
113
@@ -119,10 +123,11 @@ func RunDHTConnectivity(conf testutil.LatencyConfig, numPeers int) error {
123 lanPeer.Peerstore.AddAddr(lanPeer.Identity, lanAddr, peerstore.PermanentAddrTTL)
124 for _, p := range lanPeers {
125 mn.LinkPeers(p.Identity, lanPeer.Identity)
122 - mn.ConnectPeers(p.Identity, lanPeer.Identity)
126 + lanPeer.PeerHost.Connect(connectionContext, p.Peerstore.PeerInfo(p.Identity))
127 }
128 lanPeers = append(lanPeers, lanPeer)
129 }
130 + connCtxCancel()
131
132 // Add interfaces / addresses to test peer.
133 wanAddr := makeAddr(0, true)
@@ -136,24 +141,48 @@ func RunDHTConnectivity(conf testutil.LatencyConfig, numPeers int) error {
141 return err
142 }
143 }
139 - _, err = mn.ConnectPeers(testPeer.Identity, lanPeers[0].Identity)
144 + err = testPeer.PeerHost.Connect(ctx, lanPeers[0].Peerstore.PeerInfo(lanPeers[0].Identity))
145 if err != nil {
146 return err
147 }
148
144 - err, done := <-testPeer.DHT.LAN.RefreshRoutingTable()
145 - if err != nil || !done {
146 - if !done {
147 - err = fmt.Errorf("expected refresh routing table to close")
149 + startupCtx, startupCancel := context.WithTimeout(ctx, time.Second*15)
150 + testPeer.DHT.Bootstrap(startupCtx)
151 +StartupWait:
152 + for {
153 + select {
154 + case err, done := <-testPeer.DHT.LAN.RefreshRoutingTable():
155 + if err.Error() == kbucket.ErrLookupFailure.Error() ||
156 + testPeer.DHT.LAN.RoutingTable() == nil ||
157 + testPeer.DHT.LAN.RoutingTable().Size() == 0 {
158 + time.Sleep(100 * time.Millisecond)
159 + continue
160 + }
161 + if err != nil || !done {
162 + if !done {
163 + err = fmt.Errorf("expected refresh routing table to close")
164 + }
165 + fmt.Fprintf(os.Stderr, "how odd. that was lookupfailure.\n")
166 + startupCancel()
167 + return err
168 + }
169 + break StartupWait
170 + case <-startupCtx.Done():
171 + startupCancel()
172 + return fmt.Errorf("expected faster dht bootstrap")
173 }
149 - return err
174 }
175 + startupCancel()
176
177 + fmt.Fprintf(os.Stderr, "finding provider\n")
178 // choose a lan peer and validate lan DHT is functioning.
179 i := rand.Intn(len(lanPeers))
180 if testPeer.PeerHost.Network().Connectedness(lanPeers[i].Identity) == corenet.Connected {
155 - testPeer.PeerHost.Network().ClosePeer(lanPeers[i].Identity)
156 - testPeer.PeerHost.Peerstore().ClearAddrs(lanPeers[i].Identity)
181 + i = (i + 1) % len(lanPeers)
182 + if testPeer.PeerHost.Network().Connectedness(lanPeers[i].Identity) == corenet.Connected {
183 + testPeer.PeerHost.Network().ClosePeer(lanPeers[i].Identity)
184 + testPeer.PeerHost.Peerstore().ClearAddrs(lanPeers[i].Identity)
185 + }
186 }
187 // That peer will provide a new CID, and we'll validate the test node can find it.
188 provideCid := cid.NewCidV1(cid.Raw, []byte("Lan Provide Record"))
@@ -183,7 +212,7 @@ func RunDHTConnectivity(conf testutil.LatencyConfig, numPeers int) error {
212 return err
213 }
214
186 - err, done = <-testPeer.DHT.WAN.RefreshRoutingTable()
215 + err, done := <-testPeer.DHT.WAN.RefreshRoutingTable()
216 if err != nil || !done {
217 if !done {
218 err = fmt.Errorf("expected refresh routing table to close")