refactor(dht/pb) move proto to pb package
Brian Tiger Chow committed
Oct 25, 2014 at 04:13 UTC
29457214cba9c995eace1fecf63247b0c1b3327c
9 files changed
+86
-83
routing/dht/Makefile
deleted
-11
@@ -1,11 +0,0 @@
1
-
2
-PB = $(wildcard *.proto)
3
-GO = $(PB:.proto=.pb.go)
4
-
5
-all: $(GO)
6
-
7
-%.pb.go: %.proto
8
- protoc --gogo_out=. --proto_path=../../../../:/usr/local/opt/protobuf/include:. $<
9
-
10
-clean:
11
- rm *.pb.go
routing/dht/dht.go
+21
-20
@@ -11,6 +11,7 @@ import (
11
inet "github.com/jbenet/go-ipfs/net"
12
msg "github.com/jbenet/go-ipfs/net/message"
13
peer "github.com/jbenet/go-ipfs/peer"
14
+ pb "github.com/jbenet/go-ipfs/routing/dht/pb"
15
kb "github.com/jbenet/go-ipfs/routing/kbucket"
16
u "github.com/jbenet/go-ipfs/util"
17
@@ -128,7 +129,7 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
129
}
130
131
// deserialize msg
131
- pmes := new(Message)
132
+ pmes := new(pb.Message)
133
err := proto.Unmarshal(mData, pmes)
134
if err != nil {
135
log.Error("Error unmarshaling data")
@@ -140,7 +141,7 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
141
142
// Print out diagnostic
143
log.Debugf("%s got message type: '%s' from %s",
143
- dht.self, Message_MessageType_name[int32(pmes.GetType())], mPeer)
144
+ dht.self, pb.Message_MessageType_name[int32(pmes.GetType())], mPeer)
145
146
// get handler for this msg type.
147
handler := dht.handlerForMsgType(pmes.GetType())
@@ -174,7 +175,7 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
175
176
// sendRequest sends out a request using dht.sender, but also makes sure to
177
// measure the RTT for latency measurements.
177
-func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *Message) (*Message, error) {
178
+func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
179
180
mes, err := msg.FromObject(p, pmes)
181
if err != nil {
@@ -185,7 +186,7 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *Message)
186
187
// Print out diagnostic
188
log.Debugf("Sent message type: '%s' to %s",
188
- Message_MessageType_name[int32(pmes.GetType())], p)
189
+ pb.Message_MessageType_name[int32(pmes.GetType())], p)
190
191
rmes, err := dht.sender.SendRequest(ctx, mes)
192
if err != nil {
@@ -198,7 +199,7 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *Message)
199
rtt := time.Since(start)
200
rmes.Peer().SetLatency(rtt)
201
201
- rpmes := new(Message)
202
+ rpmes := new(pb.Message)
203
if err := proto.Unmarshal(rmes.Data(), rpmes); err != nil {
204
return nil, err
205
}
@@ -210,7 +211,7 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *Message)
211
func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.Peer,
212
key string, value []byte) error {
213
213
- pmes := newMessage(Message_PUT_VALUE, string(key), 0)
214
+ pmes := pb.NewMessage(pb.Message_PUT_VALUE, string(key), 0)
215
pmes.Value = value
216
rpmes, err := dht.sendRequest(ctx, p, pmes)
217
if err != nil {
@@ -225,10 +226,10 @@ func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.Peer,
226
227
func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) error {
228
228
- pmes := newMessage(Message_ADD_PROVIDER, string(key), 0)
229
+ pmes := pb.NewMessage(pb.Message_ADD_PROVIDER, string(key), 0)
230
231
// add self as the provider
231
- pmes.ProviderPeers = peersToPBPeers([]peer.Peer{dht.self})
232
+ pmes.ProviderPeers = pb.PeersToPBPeers([]peer.Peer{dht.self})
233
234
rpmes, err := dht.sendRequest(ctx, p, pmes)
235
if err != nil {
@@ -290,9 +291,9 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
291
292
// getValueSingle simply performs the get value RPC with the given parameters
293
func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.Peer,
293
- key u.Key, level int) (*Message, error) {
294
+ key u.Key, level int) (*pb.Message, error) {
295
295
- pmes := newMessage(Message_GET_VALUE, string(key), level)
296
+ pmes := pb.NewMessage(pb.Message_GET_VALUE, string(key), level)
297
return dht.sendRequest(ctx, p, pmes)
298
}
299
@@ -301,7 +302,7 @@ func (dht *IpfsDHT) getValueSingle(ctx context.Context, p peer.Peer,
302
// one to get the value from? Or just connect to one at a time until we get a
303
// successful connection and request the value from it?
304
func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
304
- peerlist []*Message_Peer, level int) ([]byte, error) {
305
+ peerlist []*pb.Message_Peer, level int) ([]byte, error) {
306
307
for _, pinfo := range peerlist {
308
p, err := dht.ensureConnectedToPeer(pinfo)
@@ -379,17 +380,17 @@ func (dht *IpfsDHT) FindLocal(id peer.ID) (peer.Peer, *kb.RoutingTable) {
380
return nil, nil
381
}
382
382
-func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID, level int) (*Message, error) {
383
- pmes := newMessage(Message_FIND_NODE, string(id), level)
383
+func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID, level int) (*pb.Message, error) {
384
+ pmes := pb.NewMessage(pb.Message_FIND_NODE, string(id), level)
385
return dht.sendRequest(ctx, p, pmes)
386
}
387
387
-func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.Key, level int) (*Message, error) {
388
- pmes := newMessage(Message_GET_PROVIDERS, string(key), level)
388
+func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.Peer, key u.Key, level int) (*pb.Message, error) {
389
+ pmes := pb.NewMessage(pb.Message_GET_PROVIDERS, string(key), level)
390
return dht.sendRequest(ctx, p, pmes)
391
}
392
392
-func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []peer.Peer {
393
+func (dht *IpfsDHT) addProviders(key u.Key, peers []*pb.Message_Peer) []peer.Peer {
394
var provArr []peer.Peer
395
for _, prov := range peers {
396
p, err := dht.peerFromInfo(prov)
@@ -413,7 +414,7 @@ func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []peer.Peer {
414
}
415
416
// nearestPeersToQuery returns the routing tables closest peers.
416
-func (dht *IpfsDHT) nearestPeersToQuery(pmes *Message, count int) []peer.Peer {
417
+func (dht *IpfsDHT) nearestPeersToQuery(pmes *pb.Message, count int) []peer.Peer {
418
level := pmes.GetClusterLevel()
419
cluster := dht.routingTables[level]
420
@@ -423,7 +424,7 @@ func (dht *IpfsDHT) nearestPeersToQuery(pmes *Message, count int) []peer.Peer {
424
}
425
426
// betterPeerToQuery returns nearestPeersToQuery, but iff closer than self.
426
-func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []peer.Peer {
427
+func (dht *IpfsDHT) betterPeersToQuery(pmes *pb.Message, count int) []peer.Peer {
428
closer := dht.nearestPeersToQuery(pmes, count)
429
430
// no node? nil
@@ -462,7 +463,7 @@ func (dht *IpfsDHT) getPeer(id peer.ID) (peer.Peer, error) {
463
return p, nil
464
}
465
465
-func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (peer.Peer, error) {
466
+func (dht *IpfsDHT) peerFromInfo(pbp *pb.Message_Peer) (peer.Peer, error) {
467
468
id := peer.ID(pbp.GetId())
469
@@ -485,7 +486,7 @@ func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (peer.Peer, error) {
486
return p, nil
487
}
488
488
-func (dht *IpfsDHT) ensureConnectedToPeer(pbp *Message_Peer) (peer.Peer, error) {
489
+func (dht *IpfsDHT) ensureConnectedToPeer(pbp *pb.Message_Peer) (peer.Peer, error) {
490
p, err := dht.peerFromInfo(pbp)
491
if err != nil {
492
return nil, err
routing/dht/ext_test.go
+14
-13
@@ -12,6 +12,7 @@ import (
12
msg "github.com/jbenet/go-ipfs/net/message"
13
mux "github.com/jbenet/go-ipfs/net/mux"
14
peer "github.com/jbenet/go-ipfs/peer"
15
+ pb "github.com/jbenet/go-ipfs/routing/dht/pb"
16
u "github.com/jbenet/go-ipfs/util"
17
18
"time"
@@ -127,13 +128,13 @@ func TestGetFailures(t *testing.T) {
128
// u.POut("NotFound Test\n")
129
// Reply with failures to every message
130
fs.AddHandler(func(mes msg.NetMessage) msg.NetMessage {
130
- pmes := new(Message)
131
+ pmes := new(pb.Message)
132
err := proto.Unmarshal(mes.Data(), pmes)
133
if err != nil {
134
t.Fatal(err)
135
}
136
136
- resp := &Message{
137
+ resp := &pb.Message{
138
Type: pmes.Type,
139
}
140
m, err := msg.FromObject(mes.Peer(), resp)
@@ -153,9 +154,9 @@ func TestGetFailures(t *testing.T) {
154
155
fs.handlers = nil
156
// Now we test this DHT's handleGetValue failure
156
- typ := Message_GET_VALUE
157
+ typ := pb.Message_GET_VALUE
158
str := "hello"
158
- req := Message{
159
+ req := pb.Message{
160
Type: &typ,
161
Key: &str,
162
Value: []byte{0},
@@ -169,7 +170,7 @@ func TestGetFailures(t *testing.T) {
170
171
mes = d.HandleMessage(ctx, mes)
172
172
- pmes := new(Message)
173
+ pmes := new(pb.Message)
174
err = proto.Unmarshal(mes.Data(), pmes)
175
if err != nil {
176
t.Fatal(err)
@@ -215,21 +216,21 @@ func TestNotFound(t *testing.T) {
216
217
// Reply with random peers to every message
218
fs.AddHandler(func(mes msg.NetMessage) msg.NetMessage {
218
- pmes := new(Message)
219
+ pmes := new(pb.Message)
220
err := proto.Unmarshal(mes.Data(), pmes)
221
if err != nil {
222
t.Fatal(err)
223
}
224
225
switch pmes.GetType() {
225
- case Message_GET_VALUE:
226
- resp := &Message{Type: pmes.Type}
226
+ case pb.Message_GET_VALUE:
227
+ resp := &pb.Message{Type: pmes.Type}
228
229
peers := []peer.Peer{}
230
for i := 0; i < 7; i++ {
231
peers = append(peers, _randPeer())
232
}
232
- resp.CloserPeers = peersToPBPeers(peers)
233
+ resp.CloserPeers = pb.PeersToPBPeers(peers)
234
mes, err := msg.FromObject(mes.Peer(), resp)
235
if err != nil {
236
t.Error(err)
@@ -282,17 +283,17 @@ func TestLessThanKResponses(t *testing.T) {
283
284
// Reply with random peers to every message
285
fs.AddHandler(func(mes msg.NetMessage) msg.NetMessage {
285
- pmes := new(Message)
286
+ pmes := new(pb.Message)
287
err := proto.Unmarshal(mes.Data(), pmes)
288
if err != nil {
289
t.Fatal(err)
290
}
291
292
switch pmes.GetType() {
292
- case Message_GET_VALUE:
293
- resp := &Message{
293
+ case pb.Message_GET_VALUE:
294
+ resp := &pb.Message{
295
Type: pmes.Type,
295
- CloserPeers: peersToPBPeers([]peer.Peer{other}),
296
+ CloserPeers: pb.PeersToPBPeers([]peer.Peer{other}),
297
}
298
299
mes, err := msg.FromObject(mes.Peer(), resp)
routing/dht/handlers.go
+23
-22
@@ -6,6 +6,7 @@ import (
6
"time"
7
8
peer "github.com/jbenet/go-ipfs/peer"
9
+ pb "github.com/jbenet/go-ipfs/routing/dht/pb"
10
u "github.com/jbenet/go-ipfs/util"
11
12
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
@@ -14,32 +15,32 @@ import (
15
var CloserPeerCount = 4
16
17
// dhthandler specifies the signature of functions that handle DHT messages.
17
-type dhtHandler func(peer.Peer, *Message) (*Message, error)
18
+type dhtHandler func(peer.Peer, *pb.Message) (*pb.Message, error)
19
19
-func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
20
+func (dht *IpfsDHT) handlerForMsgType(t pb.Message_MessageType) dhtHandler {
21
switch t {
21
- case Message_GET_VALUE:
22
+ case pb.Message_GET_VALUE:
23
return dht.handleGetValue
23
- case Message_PUT_VALUE:
24
+ case pb.Message_PUT_VALUE:
25
return dht.handlePutValue
25
- case Message_FIND_NODE:
26
+ case pb.Message_FIND_NODE:
27
return dht.handleFindPeer
27
- case Message_ADD_PROVIDER:
28
+ case pb.Message_ADD_PROVIDER:
29
return dht.handleAddProvider
29
- case Message_GET_PROVIDERS:
30
+ case pb.Message_GET_PROVIDERS:
31
return dht.handleGetProviders
31
- case Message_PING:
32
+ case pb.Message_PING:
33
return dht.handlePing
34
default:
35
return nil
36
}
37
}
38
38
-func (dht *IpfsDHT) handleGetValue(p peer.Peer, pmes *Message) (*Message, error) {
39
+func (dht *IpfsDHT) handleGetValue(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
40
log.Debugf("%s handleGetValue for key: %s\n", dht.self, pmes.GetKey())
41
42
// setup response
42
- resp := newMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
43
+ resp := pb.NewMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
44
45
// first, is the key even a key?
46
key := pmes.GetKey()
@@ -77,7 +78,7 @@ func (dht *IpfsDHT) handleGetValue(p peer.Peer, pmes *Message) (*Message, error)
78
provs := dht.providers.GetProviders(u.Key(pmes.GetKey()))
79
if len(provs) > 0 {
80
log.Debugf("handleGetValue returning %d provider[s]", len(provs))
80
- resp.ProviderPeers = peersToPBPeers(provs)
81
+ resp.ProviderPeers = pb.PeersToPBPeers(provs)
82
}
83
84
// Find closest peer on given cluster to desired key and reply with that info
@@ -89,14 +90,14 @@ func (dht *IpfsDHT) handleGetValue(p peer.Peer, pmes *Message) (*Message, error)
90
log.Critical("no addresses on peer being sent!")
91
}
92
}
92
- resp.CloserPeers = peersToPBPeers(closer)
93
+ resp.CloserPeers = pb.PeersToPBPeers(closer)
94
}
95
96
return resp, nil
97
}
98
99
// Store a value in this peer local storage
99
-func (dht *IpfsDHT) handlePutValue(p peer.Peer, pmes *Message) (*Message, error) {
100
+func (dht *IpfsDHT) handlePutValue(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
101
dht.dslock.Lock()
102
defer dht.dslock.Unlock()
103
dskey := u.Key(pmes.GetKey()).DsKey()
@@ -105,13 +106,13 @@ func (dht *IpfsDHT) handlePutValue(p peer.Peer, pmes *Message) (*Message, error)
106
return pmes, err
107
}
108
108
-func (dht *IpfsDHT) handlePing(p peer.Peer, pmes *Message) (*Message, error) {
109
+func (dht *IpfsDHT) handlePing(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
110
log.Debugf("%s Responding to ping from %s!\n", dht.self, p)
111
return pmes, nil
112
}
113
113
-func (dht *IpfsDHT) handleFindPeer(p peer.Peer, pmes *Message) (*Message, error) {
114
- resp := newMessage(pmes.GetType(), "", pmes.GetClusterLevel())
114
+func (dht *IpfsDHT) handleFindPeer(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
115
+ resp := pb.NewMessage(pmes.GetType(), "", pmes.GetClusterLevel())
116
var closest []peer.Peer
117
118
// if looking for self... special case where we send it on CloserPeers.
@@ -136,12 +137,12 @@ func (dht *IpfsDHT) handleFindPeer(p peer.Peer, pmes *Message) (*Message, error)
137
for _, p := range withAddresses {
138
log.Debugf("handleFindPeer: sending back '%s'", p)
139
}
139
- resp.CloserPeers = peersToPBPeers(withAddresses)
140
+ resp.CloserPeers = pb.PeersToPBPeers(withAddresses)
141
return resp, nil
142
}
143
143
-func (dht *IpfsDHT) handleGetProviders(p peer.Peer, pmes *Message) (*Message, error) {
144
- resp := newMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
144
+func (dht *IpfsDHT) handleGetProviders(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
145
+ resp := pb.NewMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
146
147
// check if we have this value, to add ourselves as provider.
148
log.Debugf("handling GetProviders: '%s'", pmes.GetKey())
@@ -160,13 +161,13 @@ func (dht *IpfsDHT) handleGetProviders(p peer.Peer, pmes *Message) (*Message, er
161
162
// if we've got providers, send thos those.
163
if providers != nil && len(providers) > 0 {
163
- resp.ProviderPeers = peersToPBPeers(providers)
164
+ resp.ProviderPeers = pb.PeersToPBPeers(providers)
165
}
166
167
// Also send closer peers.
168
closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
169
if closer != nil {
169
- resp.CloserPeers = peersToPBPeers(closer)
170
+ resp.CloserPeers = pb.PeersToPBPeers(closer)
171
}
172
173
return resp, nil
@@ -177,7 +178,7 @@ type providerInfo struct {
178
Value peer.Peer
179
}
180
180
-func (dht *IpfsDHT) handleAddProvider(p peer.Peer, pmes *Message) (*Message, error) {
181
+func (dht *IpfsDHT) handleAddProvider(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
182
key := u.Key(pmes.GetKey())
183
184
log.Debugf("%s adding %s as a provider for '%s'\n", dht.self, p, peer.ID(key))
routing/dht/pb/Makefile
new
+11
@@ -0,0 +1,11 @@
1
+PB = $(wildcard *.proto)
2
+GO = $(PB:.proto=.pb.go)
3
+
4
+all: $(GO)
5
+
6
+%.pb.go: %.proto
7
+ protoc --gogo_out=. --proto_path=../../../../../../:/usr/local/opt/protobuf/include:. $<
8
+
9
+clean:
10
+ rm -f *.pb.go
11
+ rm -f *.go
routing/dht/pb/dht.pb.go
renamed
+8
-8
@@ -1,19 +1,19 @@
1
-// Code generated by protoc-gen-go.
2
-// source: messages.proto
1
+// Code generated by protoc-gen-gogo.
2
+// source: dht.proto
3
// DO NOT EDIT!
4
5
/*
6
-Package dht is a generated protocol buffer package.
6
+Package dht_pb is a generated protocol buffer package.
7
8
It is generated from these files:
9
- messages.proto
9
+ dht.proto
10
11
It has these top-level messages:
12
Message
13
*/
14
-package dht
14
+package dht_pb
15
16
-import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
16
+import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/proto"
17
import math "math"
18
19
// Reference imports to suppress errors if they are not otherwise used.
@@ -67,7 +67,7 @@ func (x *Message_MessageType) UnmarshalJSON(data []byte) error {
67
68
type Message struct {
69
// defines what type of message it is.
70
- Type *Message_MessageType `protobuf:"varint,1,opt,name=type,enum=dht.Message_MessageType" json:"type,omitempty"`
70
+ Type *Message_MessageType `protobuf:"varint,1,opt,name=type,enum=dht.pb.Message_MessageType" json:"type,omitempty"`
71
// defines what coral cluster level this query/response belongs to.
72
ClusterLevelRaw *int32 `protobuf:"varint,10,opt,name=clusterLevelRaw" json:"clusterLevelRaw,omitempty"`
73
// Used to specify the key associated with this message.
@@ -156,5 +156,5 @@ func (m *Message_Peer) GetAddr() string {
156
}
157
158
func init() {
159
- proto.RegisterEnum("dht.Message_MessageType", Message_MessageType_name, Message_MessageType_value)
159
+ proto.RegisterEnum("dht.pb.Message_MessageType", Message_MessageType_name, Message_MessageType_value)
160
}
routing/dht/pb/dht.proto
renamed
+1
-1
@@ -1,4 +1,4 @@
1
-package dht;
1
+package dht.pb;
2
3
//run `protoc --go_out=. *.proto` to generate
4
routing/dht/pb/message.go
renamed
+4
-5
@@ -1,4 +1,4 @@
1
-package dht
1
+package dht_pb
2
3
import (
4
"errors"
@@ -8,7 +8,7 @@ import (
8
peer "github.com/jbenet/go-ipfs/peer"
9
)
10
11
-func newMessage(typ Message_MessageType, key string, level int) *Message {
11
+func NewMessage(typ Message_MessageType, key string, level int) *Message {
12
m := &Message{
13
Type: &typ,
14
Key: &key,
@@ -31,7 +31,7 @@ func peerToPBPeer(p peer.Peer) *Message_Peer {
31
return pbp
32
}
33
34
-func peersToPBPeers(peers []peer.Peer) []*Message_Peer {
34
+func PeersToPBPeers(peers []peer.Peer) []*Message_Peer {
35
pbpeers := make([]*Message_Peer, len(peers))
36
for i, p := range peers {
37
pbpeers[i] = peerToPBPeer(p)
@@ -53,8 +53,7 @@ func (m *Message_Peer) Address() (ma.Multiaddr, error) {
53
func (m *Message) GetClusterLevel() int {
54
level := m.GetClusterLevelRaw() - 1
55
if level < 0 {
56
- log.Debug("GetClusterLevel: no routing level specified, assuming 0")
57
- level = 0
56
+ return 0
57
}
58
return int(level)
59
}
routing/dht/routing.go
+4
-3
@@ -6,6 +6,7 @@ import (
6
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
8
peer "github.com/jbenet/go-ipfs/peer"
9
+ pb "github.com/jbenet/go-ipfs/routing/dht/pb"
10
kb "github.com/jbenet/go-ipfs/routing/kbucket"
11
u "github.com/jbenet/go-ipfs/util"
12
)
@@ -152,10 +153,10 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
153
return peerOut
154
}
155
155
-func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet, count int, out chan peer.Peer) {
156
+func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*pb.Message_Peer, ps *peerSet, count int, out chan peer.Peer) {
157
done := make(chan struct{})
158
for _, pbp := range peers {
158
- go func(mp *Message_Peer) {
159
+ go func(mp *pb.Message_Peer) {
160
defer func() { done <- struct{}{} }()
161
// construct new peer
162
p, err := dht.ensureConnectedToPeer(mp)
@@ -258,7 +259,7 @@ func (dht *IpfsDHT) Ping(ctx context.Context, p peer.Peer) error {
259
// Thoughts: maybe this should accept an ID and do a peer lookup?
260
log.Infof("ping %s start", p)
261
261
- pmes := newMessage(Message_PING, "", 0)
262
+ pmes := pb.NewMessage(pb.Message_PING, "", 0)
263
_, err := dht.sendRequest(ctx, p, pmes)
264
log.Infof("ping %s end (err = %s)", p, err)
265
return err