@cryptotaxi247 / kubo / commits / 9eb41e723

ping + find peer

Juan Batiz-Benet committed Sep 16, 2014 at 01:09 UTC 9eb41e7237bbd66c64d7e4a8f32e9d1fca51f4e4
2 files changed +36 -38
routing/dht/Message.go
+16 -2
@@ -3,9 +3,10 @@ package dht
3 import (
4 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
5 peer "github.com/jbenet/go-ipfs/peer"
6 + u "github.com/jbenet/go-ipfs/util"
7 )
8
8 -func peerInfo(p *peer.Peer) *Message_Peer {
9 +func peerToPBPeer(p *peer.Peer) *Message_Peer {
10 pbp := new(Message_Peer)
11 if len(p.Addresses) == 0 || p.Addresses[0] == nil {
12 pbp.Addr = proto.String("")
@@ -22,11 +23,24 @@ func peerInfo(p *peer.Peer) *Message_Peer {
23 return pbp
24 }
25
26 +func peersToPBPeers(peers []*peer.Peer) []*Message_Peer {
27 + pbpeers = make([]*Message_Peer, len(peers))
28 + for i, p := range peers {
29 + pbpeers[i] = peerToPBPeer(p)
30 + }
31 + return pbpeers
32 +}
33 +
34 // GetClusterLevel gets and adjusts the cluster level on the message.
35 // a +/- 1 adjustment is needed to distinguish a valid first level (1) and
36 // default "no value" protobuf behavior (0)
37 func (m *Message) GetClusterLevel() int32 {
29 - return m.GetClusterLevelRaw() - 1
38 + level := m.GetClusterLevelRaw() - 1
39 + if level < 0 {
40 + u.PErr("handleGetValue: no routing level specified, assuming 0\n")
41 + level = 0
42 + }
43 + return level
44 }
45
46 // SetClusterLevel adjusts and sets the cluster level on the message.
routing/dht/dht.go
+20 -36
@@ -169,14 +169,14 @@ func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
169 return dht.handleGetValue
170 // case Message_PUT_VALUE:
171 // return dht.handlePutValue
172 - // case Message_FIND_NODE:
173 - // return dht.handleFindPeer
172 + case Message_FIND_NODE:
173 + return dht.handleFindPeer
174 // case Message_ADD_PROVIDER:
175 // return dht.handleAddProvider
176 // case Message_GET_PROVIDERS:
177 // return dht.handleGetProviders
178 - // case Message_PING:
179 - // return dht.handlePing
178 + case Message_PING:
179 + return dht.handlePing
180 // case Message_DIAGNOSTIC:
181 // return dht.handleDiagnostic
182 default:
@@ -240,7 +240,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
240 provs := dht.providers.GetProviders(u.Key(pmes.GetKey()))
241 if len(provs) > 0 {
242 u.DOut("handleGetValue returning %d provider[s]\n", len(provs))
243 - resp.ProviderPeers = provs
243 + resp.ProviderPeers = peersToPBPeers(provs)
244 return resp, nil
245 }
246
@@ -249,11 +249,6 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
249
250 // stored levels are > 1, to distinguish missing levels.
251 level := pmes.GetClusterLevel()
252 - if level < 0 {
253 - // TODO: maybe return an error? Defaulting isnt a good idea IMO
254 - u.PErr("handleGetValue: no routing level specified, assuming 0\n")
255 - level = 0
256 - }
252 u.DOut("handleGetValue searching level %d clusters\n", level)
253
254 ck := kb.ConvertKey(u.Key(pmes.GetKey()))
@@ -275,7 +270,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
270
271 // we got a closer peer, it seems. return it.
272 u.DOut("handleGetValue returning a closer peer: '%s'\n", closer.ID.Pretty())
278 - resp.CloserPeers = []*peer.Peer{closer}
273 + resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
274 return resp, nil
275 }
276
@@ -291,48 +286,37 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) {
286 }
287 }
288
294 -func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) {
289 +func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
290 u.DOut("[%s] Responding to ping from [%s]!\n", dht.self.ID.Pretty(), p.ID.Pretty())
296 - resp := Message{
297 - Type: pmes.GetType(),
298 - Response: true,
299 - ID: pmes.GetId(),
300 - }
301 -
302 - dht.netChan.Outgoing <- swarm.NewMessage(p, resp.ToProtobuf())
291 + return &Message{Type: pmes.Type}, nil
292 }
293
305 -func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) {
306 - resp := Message{
307 - Type: pmes.GetType(),
308 - ID: pmes.GetId(),
309 - Response: true,
310 - }
311 - defer func() {
312 - mes := swarm.NewMessage(p, resp.ToProtobuf())
313 - dht.netChan.Outgoing <- mes
314 - }()
315 - level := pmes.GetValue()[0]
294 +func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error) {
295 + resp := &Message{Type: pmes.Type}
296 +
297 + level := pmes.GetClusterLevel()
298 u.DOut("handleFindPeer: searching for '%s'\n", peer.ID(pmes.GetKey()).Pretty())
317 - closest := dht.routingTables[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
299 +
300 + ck := kb.ConvertKey(u.Key(pmes.GetKey()))
301 + closest := dht.routingTables[level].NearestPeer(ck)
302 if closest == nil {
303 u.PErr("handleFindPeer: could not find anything.\n")
320 - return
304 + return resp, nil
305 }
306
307 if len(closest.Addresses) == 0 {
308 u.PErr("handleFindPeer: no addresses for connected peer...\n")
325 - return
309 + return resp, nil
310 }
311
312 // If the found peer further away than this peer...
313 if kb.Closer(dht.self.ID, closest.ID, u.Key(pmes.GetKey())) {
330 - return
314 + return resp, nil
315 }
316
317 u.DOut("handleFindPeer: sending back '%s'\n", closest.ID.Pretty())
334 - resp.Peers = []*peer.Peer{closest}
335 - resp.Success = true
318 + resp.CloserPeers = peersToPBPeers([]*peer.Peer{closest})
319 + return resp, nil
320 }
321
322 func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) {