dht: update to use net.LocalPeer
Juan Batiz-Benet committed
Nov 21, 2014 at 08:03 UTC
d06bb6d8263239d4a4bc30cfc5bf06bcf28b60bc
4 files changed
+64
-61
routing/dht/dht.go
+14
-43
@@ -274,14 +274,11 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
274
}
275
276
// Perhaps we were given closer peers
277
- var peers []peer.Peer
278
- for _, pb := range pmes.GetCloserPeers() {
279
- pr, err := dht.peerFromInfo(pb)
277
+ peers, errs := pb.PBPeersToPeers(dht.peerstore, pmes.GetCloserPeers())
278
+ for _, err := range errs {
279
if err != nil {
280
log.Error(err)
282
- continue
281
}
284
- peers = append(peers, pr)
282
}
283
284
if len(peers) > 0 {
@@ -426,22 +423,20 @@ func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.
423
return dht.sendRequest(ctx, p, pmes)
424
}
425
429
-func (dht *IpfsDHT) addProviders(key u.Key, peers []*pb.Message_Peer) []peer.Peer {
430
- var provArr []peer.Peer
431
- for _, prov := range peers {
432
- p, err := dht.peerFromInfo(prov)
433
- if err != nil {
434
- log.Errorf("error getting peer from info: %v", err)
435
- continue
436
- }
437
-
438
- log.Debugf("%s adding provider: %s for %s", dht.self, p, key)
426
+func (dht *IpfsDHT) addProviders(key u.Key, pbps []*pb.Message_Peer) []peer.Peer {
427
+ peers, errs := pb.PBPeersToPeers(dht.peerstore, pbps)
428
+ for _, err := range errs {
429
+ log.Errorf("error converting peer: %v", err)
430
+ }
431
432
+ var provArr []peer.Peer
433
+ for _, p := range peers {
434
// Dont add outselves to the list
435
if p.ID().Equal(dht.self.ID()) {
436
continue
437
}
438
439
+ log.Debugf("%s adding provider: %s for %s", dht.self, p, key)
440
// TODO(jbenet) ensure providers is idempotent
441
dht.providers.AddProvider(key, p)
442
provArr = append(provArr, p)
@@ -500,38 +495,14 @@ func (dht *IpfsDHT) getPeer(id peer.ID) (peer.Peer, error) {
495
return p, nil
496
}
497
503
-// peerFromInfo returns a peer using info in the protobuf peer struct
504
-// to lookup or create a peer
505
-func (dht *IpfsDHT) peerFromInfo(pbp *pb.Message_Peer) (peer.Peer, error) {
506
-
507
- id := peer.ID(pbp.GetId())
508
-
509
- // bail out if it's ourselves
510
- //TODO(jbenet) not sure this should be an error _here_
511
- if id.Equal(dht.self.ID()) {
512
- return nil, errors.New("found self")
513
- }
514
-
515
- p, err := dht.getPeer(id)
516
- if err != nil {
517
- return nil, err
518
- }
519
-
520
- // add addresses we've just discovered
521
- maddrs, err := pbp.Addresses()
498
+func (dht *IpfsDHT) ensureConnectedToPeer(ctx context.Context, pbp *pb.Message_Peer) (peer.Peer, error) {
499
+ p, err := pb.PBPeerToPeer(dht.peerstore, pbp)
500
if err != nil {
501
return nil, err
502
}
525
- for _, maddr := range maddrs {
526
- p.AddAddress(maddr)
527
- }
528
- return p, nil
529
-}
503
531
-func (dht *IpfsDHT) ensureConnectedToPeer(ctx context.Context, pbp *pb.Message_Peer) (peer.Peer, error) {
532
- p, err := dht.peerFromInfo(pbp)
533
- if err != nil {
534
- return nil, err
504
+ if dht.dialer.LocalPeer().ID().Equal(p.ID()) {
505
+ return nil, errors.New("attempting to ensure connection to self")
506
}
507
508
// dial connection
routing/dht/pb/message.go
+37
@@ -2,6 +2,7 @@ package dht_pb
2
3
import (
4
"errors"
5
+ "fmt"
6
7
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
8
@@ -32,6 +33,24 @@ func peerToPBPeer(p peer.Peer) *Message_Peer {
33
return pbp
34
}
35
36
+// PBPeerToPeer turns a *Message_Peer into its peer.Peer counterpart
37
+func PBPeerToPeer(ps peer.Peerstore, pbp *Message_Peer) (peer.Peer, error) {
38
+ p, err := ps.FindOrCreate(peer.ID(pbp.GetId()))
39
+ if err != nil {
40
+ return nil, fmt.Errorf("Failed to get peer from peerstore: %s", err)
41
+ }
42
+
43
+ // add addresses
44
+ maddrs, err := pbp.Addresses()
45
+ if err != nil {
46
+ return nil, fmt.Errorf("Received peer with bad or missing addresses: %s", pbp.Addrs)
47
+ }
48
+ for _, maddr := range maddrs {
49
+ p.AddAddress(maddr)
50
+ }
51
+ return p, nil
52
+}
53
+
54
// RawPeersToPBPeers converts a slice of Peers into a slice of *Message_Peers,
55
// ready to go out on the wire.
56
func RawPeersToPBPeers(peers []peer.Peer) []*Message_Peer {
@@ -55,6 +74,24 @@ func PeersToPBPeers(d inet.Dialer, peers []peer.Peer) []*Message_Peer {
74
return pbps
75
}
76
77
+// PBPeersToPeers converts given []*Message_Peer into a set of []peer.Peer
78
+// Returns two slices, one of peers, and one of errors. The slice of peers
79
+// will ONLY contain successfully converted peers. The slice of errors contains
80
+// whether each input Message_Peer was successfully converted.
81
+func PBPeersToPeers(ps peer.Peerstore, pbps []*Message_Peer) ([]peer.Peer, []error) {
82
+ errs := make([]error, len(pbps))
83
+ peers := make([]peer.Peer, 0, len(pbps))
84
+ for i, pbp := range pbps {
85
+ p, err := PBPeerToPeer(ps, pbp)
86
+ if err != nil {
87
+ errs[i] = err
88
+ } else {
89
+ peers = append(peers, p)
90
+ }
91
+ }
92
+ return peers, errs
93
+}
94
+
95
// Addresses returns a multiaddr associated with the Message_Peer entry
96
func (m *Message_Peer) Addresses() ([]ma.Multiaddr, error) {
97
if m == nil {
routing/dht/query.go
+6
-1
@@ -161,7 +161,12 @@ func (r *dhtQueryRunner) addPeerToQuery(next peer.Peer, benchmark peer.Peer) {
161
return
162
}
163
164
- // if new peer further away than whom we got it from, bother (loops)
164
+ // if new peer is ourselves...
165
+ if next.ID().Equal(r.query.dialer.LocalPeer().ID()) {
166
+ return
167
+ }
168
+
169
+ // if new peer further away than whom we got it from, don't bother (loops)
170
if benchmark != nil && kb.Closer(benchmark.ID(), next.ID(), r.query.key) {
171
return
172
}
routing/dht/routing.go
+7
-17
@@ -234,31 +234,21 @@ func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.Peer, error)
234
}
235
236
closer := pmes.GetCloserPeers()
237
- var clpeers []peer.Peer
238
- for _, pbp := range closer {
239
- np, err := dht.getPeer(peer.ID(pbp.GetId()))
237
+ clpeers, errs := pb.PBPeersToPeers(dht.peerstore, closer)
238
+ for _, err := range errs {
239
if err != nil {
241
- log.Warningf("Received invalid peer from query: %v", err)
242
- continue
243
- }
244
-
245
- // add addresses
246
- maddrs, err := pbp.Addresses()
247
- if err != nil {
248
- log.Warning("Received peer with bad or missing addresses: %s", pbp.Addrs)
249
- continue
250
- }
251
- for _, maddr := range maddrs {
252
- np.AddAddress(maddr)
240
+ log.Warning(err)
241
}
242
+ }
243
255
- if pbp.GetId() == string(id) {
244
+ // see it we got the peer here
245
+ for _, np := range clpeers {
246
+ if string(np.ID()) == string(id) {
247
return &dhtQueryResult{
248
peer: np,
249
success: true,
250
}, nil
251
}
261
- clpeers = append(clpeers, np)
252
}
253
254
return &dhtQueryResult{closerPeers: clpeers}, nil