master
go 278 lines 7.87 KB
Raw
1 package integrationtest
2
3 import (
4 "context"
5 "encoding/binary"
6 "fmt"
7 "math"
8 "math/rand"
9 "net"
10 "testing"
11 "time"
12
13 "github.com/ipfs/go-cid"
14 "github.com/ipfs/kubo/core"
15 mock "github.com/ipfs/kubo/core/mock"
16 libp2p2 "github.com/ipfs/kubo/core/node/libp2p"
17
18 testutil "github.com/libp2p/go-libp2p-testing/net"
19 corenet "github.com/libp2p/go-libp2p/core/network"
20 mocknet "github.com/libp2p/go-libp2p/p2p/net/mock"
21
22 ma "github.com/multiformats/go-multiaddr"
23 )
24
25 func TestDHTConnectivityFast(t *testing.T) {
26 conf := testutil.LatencyConfig{
27 NetworkLatency: 0,
28 RoutingLatency: 0,
29 BlockstoreLatency: 0,
30 }
31 if err := RunDHTConnectivity(conf, 5); err != nil {
32 t.Fatal(err)
33 }
34 }
35
36 func TestDHTConnectivitySlowNetwork(t *testing.T) {
37 SkipUnlessEpic(t)
38 conf := testutil.LatencyConfig{NetworkLatency: 400 * time.Millisecond}
39 if err := RunDHTConnectivity(conf, 5); err != nil {
40 t.Fatal(err)
41 }
42 }
43
44 func TestDHTConnectivitySlowRouting(t *testing.T) {
45 SkipUnlessEpic(t)
46 conf := testutil.LatencyConfig{RoutingLatency: 400 * time.Millisecond}
47 if err := RunDHTConnectivity(conf, 5); err != nil {
48 t.Fatal(err)
49 }
50 }
51
52 // wan prefix must have a real corresponding ASN for the peer diversity filter to work.
53 var (
54 wanPrefix = net.ParseIP("2001:218:3004::")
55 lanPrefix = net.ParseIP("fe80::")
56 )
57
58 func makeAddr(n uint32, wan bool) ma.Multiaddr {
59 var ip net.IP
60 if wan {
61 ip = append(net.IP{}, wanPrefix...)
62 } else {
63 ip = append(net.IP{}, lanPrefix...)
64 }
65
66 binary.LittleEndian.PutUint32(ip[12:], n)
67 addr, _ := ma.NewMultiaddr(fmt.Sprintf("/ip6/%s/tcp/4242", ip))
68 return addr
69 }
70
71 func RunDHTConnectivity(conf testutil.LatencyConfig, numPeers int) error {
72 ctx, cancel := context.WithCancel(context.Background())
73 defer cancel()
74
75 // create network
76 mn := mocknet.New()
77 mn.SetLinkDefaults(mocknet.LinkOptions{
78 Latency: conf.NetworkLatency,
79 Bandwidth: math.MaxInt32,
80 })
81
82 testPeer, err := core.NewNode(ctx, &core.BuildCfg{
83 Online: true,
84 Host: mock.MockHostOption(mn),
85 })
86 if err != nil {
87 return err
88 }
89 defer testPeer.Close()
90
91 wanPeers := []*core.IpfsNode{}
92 lanPeers := []*core.IpfsNode{}
93
94 connectionContext, connCtxCancel := context.WithTimeout(ctx, 15*time.Second)
95 defer connCtxCancel()
96 for i := range numPeers {
97 wanPeer, err := core.NewNode(ctx, &core.BuildCfg{
98 Online: true,
99 Routing: libp2p2.DHTServerOption,
100 Host: mock.MockHostOption(mn),
101 })
102 if err != nil {
103 return err
104 }
105 defer wanPeer.Close()
106 wanAddr := makeAddr(uint32(i), true)
107 _ = wanPeer.PeerHost.Network().Listen(wanAddr)
108 for _, p := range wanPeers {
109 _, _ = mn.LinkPeers(p.Identity, wanPeer.Identity)
110 _ = wanPeer.PeerHost.Connect(connectionContext, p.Peerstore.PeerInfo(p.Identity))
111 }
112 wanPeers = append(wanPeers, wanPeer)
113
114 lanPeer, err := core.NewNode(ctx, &core.BuildCfg{
115 Online: true,
116 Host: mock.MockHostOption(mn),
117 })
118 if err != nil {
119 return err
120 }
121 defer lanPeer.Close()
122 lanAddr := makeAddr(uint32(i), false)
123 _ = lanPeer.PeerHost.Network().Listen(lanAddr)
124 for _, p := range lanPeers {
125 _, _ = mn.LinkPeers(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)
134 _ = testPeer.PeerHost.Network().Listen(wanAddr)
135 lanAddr := makeAddr(0, false)
136 _ = testPeer.PeerHost.Network().Listen(lanAddr)
137 // The test peer is connected to one lan peer.
138 for _, p := range lanPeers {
139 if _, err := mn.LinkPeers(testPeer.Identity, p.Identity); err != nil {
140 return err
141 }
142 }
143 err = testPeer.PeerHost.Connect(ctx, lanPeers[0].Peerstore.PeerInfo(lanPeers[0].Identity))
144 if err != nil {
145 return err
146 }
147
148 startupCtx, startupCancel := context.WithTimeout(ctx, time.Second*60)
149 StartupWait:
150 for {
151 select {
152 case err := <-testPeer.DHT.LAN.RefreshRoutingTable():
153 if err != nil {
154 fmt.Printf("Error refreshing routing table: %v\n", err)
155 }
156 if testPeer.DHT.LAN.RoutingTable() == nil ||
157 testPeer.DHT.LAN.RoutingTable().Size() == 0 ||
158 err != nil {
159 time.Sleep(100 * time.Millisecond)
160 continue
161 }
162 break StartupWait
163 case <-startupCtx.Done():
164 startupCancel()
165 return fmt.Errorf("expected faster dht bootstrap")
166 }
167 }
168 startupCancel()
169
170 // choose a lan peer and validate lan DHT is functioning.
171 i := rand.Intn(len(lanPeers))
172 if testPeer.PeerHost.Network().Connectedness(lanPeers[i].Identity) == corenet.Connected {
173 i = (i + 1) % len(lanPeers)
174 if testPeer.PeerHost.Network().Connectedness(lanPeers[i].Identity) == corenet.Connected {
175 _ = testPeer.PeerHost.Network().ClosePeer(lanPeers[i].Identity)
176 testPeer.PeerHost.Peerstore().ClearAddrs(lanPeers[i].Identity)
177 }
178 }
179 // That peer will provide a new CID, and we'll validate the test node can find it.
180 provideCid := cid.NewCidV1(cid.Raw, []byte("Lan Provide Record"))
181 provideCtx, cancel := context.WithTimeout(ctx, time.Second)
182 defer cancel()
183 if err := lanPeers[i].DHT.Provide(provideCtx, provideCid, true); err != nil {
184 return err
185 }
186 provChan := testPeer.DHT.FindProvidersAsync(provideCtx, provideCid, 0)
187 prov, ok := <-provChan
188 if !ok || prov.ID == "" {
189 return fmt.Errorf("Expected provider. stream closed early")
190 }
191 if prov.ID != lanPeers[i].Identity {
192 return fmt.Errorf("Unexpected lan peer provided record")
193 }
194
195 // Now, connect with a wan peer.
196 for _, p := range wanPeers {
197 if _, err := mn.LinkPeers(testPeer.Identity, p.Identity); err != nil {
198 return err
199 }
200 }
201
202 err = testPeer.PeerHost.Connect(ctx, wanPeers[0].Peerstore.PeerInfo(wanPeers[0].Identity))
203 if err != nil {
204 return err
205 }
206
207 startupCtx, startupCancel = context.WithTimeout(ctx, time.Second*60)
208 WanStartupWait:
209 for {
210 select {
211 case err := <-testPeer.DHT.WAN.RefreshRoutingTable():
212 // if err != nil {
213 // fmt.Printf("Error refreshing routing table: %v\n", err)
214 // }
215 if testPeer.DHT.WAN.RoutingTable() == nil ||
216 testPeer.DHT.WAN.RoutingTable().Size() == 0 ||
217 err != nil {
218 time.Sleep(100 * time.Millisecond)
219 continue
220 }
221 break WanStartupWait
222 case <-startupCtx.Done():
223 startupCancel()
224 return fmt.Errorf("expected faster wan dht bootstrap")
225 }
226 }
227 startupCancel()
228
229 // choose a wan peer and validate wan DHT is functioning.
230 i = rand.Intn(len(wanPeers))
231 if testPeer.PeerHost.Network().Connectedness(wanPeers[i].Identity) == corenet.Connected {
232 i = (i + 1) % len(wanPeers)
233 if testPeer.PeerHost.Network().Connectedness(wanPeers[i].Identity) == corenet.Connected {
234 _ = testPeer.PeerHost.Network().ClosePeer(wanPeers[i].Identity)
235 testPeer.PeerHost.Peerstore().ClearAddrs(wanPeers[i].Identity)
236 }
237 }
238
239 // That peer will provide a new CID, and we'll validate the test node can find it.
240 wanCid := cid.NewCidV1(cid.Raw, []byte("Wan Provide Record"))
241 wanProvideCtx, cancel := context.WithTimeout(ctx, time.Second)
242 defer cancel()
243 if err := wanPeers[i].DHT.Provide(wanProvideCtx, wanCid, true); err != nil {
244 return err
245 }
246 provChan = testPeer.DHT.FindProvidersAsync(wanProvideCtx, wanCid, 0)
247 prov, ok = <-provChan
248 if !ok || prov.ID == "" {
249 return fmt.Errorf("Expected one provider, closed early")
250 }
251 if prov.ID != wanPeers[i].Identity {
252 return fmt.Errorf("Unexpected lan peer provided record")
253 }
254
255 // Finally, re-share the lan provided cid from a wan peer and expect a merged result.
256 i = rand.Intn(len(wanPeers))
257 if testPeer.PeerHost.Network().Connectedness(wanPeers[i].Identity) == corenet.Connected {
258 _ = testPeer.PeerHost.Network().ClosePeer(wanPeers[i].Identity)
259 testPeer.PeerHost.Peerstore().ClearAddrs(wanPeers[i].Identity)
260 }
261
262 provideCtx, cancel = context.WithTimeout(ctx, time.Second)
263 defer cancel()
264 if err := wanPeers[i].DHT.Provide(provideCtx, provideCid, true); err != nil {
265 return err
266 }
267 provChan = testPeer.DHT.FindProvidersAsync(provideCtx, provideCid, 0)
268 prov, ok = <-provChan
269 if !ok {
270 return fmt.Errorf("Expected two providers, got 0")
271 }
272 prov, ok = <-provChan
273 if !ok {
274 return fmt.Errorf("Expected two providers, got 1")
275 }
276
277 return nil
278 }