@cryptotaxi247 / kubo / commits / 24bfbfe37

implement timeouts on listeners for the dht and add diagnostic stuff

Jeromy committed Aug 7, 2014 at 18:06 UTC 24bfbfe37252b0f312f90428ecdebd4e6855384a
4 files changed +130 -25
routing/dht/dht.go
+55 -9
@@ -23,6 +23,8 @@ import (
23 // IpfsDHT is an implementation of Kademlia with Coral and S/Kademlia modifications.
24 // It is used to implement the base IpfsRouting module.
25 type IpfsDHT struct {
26 + // Array of routing tables for differently distanced nodes
27 + // NOTE: (currently, only a single table is used)
28 routes []*RoutingTable
29
30 network *swarm.Swarm
@@ -55,6 +57,7 @@ type IpfsDHT struct {
57 type listenInfo struct {
58 resp chan *swarm.Message
59 count int
60 + eol time.Time
61 }
62
63 // Create a new DHT object with the given peer as the 'local' host
@@ -161,14 +164,19 @@ func (dht *IpfsDHT) handleMessages() {
164 if pmes.GetResponse() {
165 dht.listenLock.RLock()
166 list, ok := dht.listeners[pmes.GetId()]
167 + dht.listenLock.RUnlock()
168 + if time.Now().After(list.eol) {
169 + dht.Unlisten(pmes.GetId())
170 + ok = false
171 + }
172 if list.count > 1 {
173 list.count--
166 - } else if list.count == 1 {
167 - delete(dht.listeners, pmes.GetId())
174 }
169 - dht.listenLock.RUnlock()
175 if ok {
176 list.resp <- mes
177 + if list.count == 1 {
178 + dht.Unlisten(pmes.GetId())
179 + }
180 } else {
181 // this is expected behaviour during a timeout
182 u.DOut("Received response with nobody listening...")
@@ -217,10 +225,35 @@ func (dht *IpfsDHT) handleMessages() {
225 dht.providers[k] = cleaned
226 }
227 dht.providerLock.Unlock()
228 + dht.listenLock.Lock()
229 + var remove []uint64
230 + now := time.Now()
231 + for k,v := range dht.listeners {
232 + if now.After(v.eol) {
233 + remove = append(remove, k)
234 + }
235 + }
236 + for _,k := range remove {
237 + delete(dht.listeners, k)
238 + }
239 + dht.listenLock.Unlock()
240 }
241 }
242 }
243
244 +func (dht *IpfsDHT) putValueToPeer(p *peer.Peer, key string, value []byte) error {
245 + pmes := pDHTMessage{
246 + Type: DHTMessage_PUT_VALUE,
247 + Key: key,
248 + Value: value,
249 + Id: GenerateMessageID(),
250 + }
251 +
252 + mes := swarm.NewMessage(p, pmes.ToProtobuf())
253 + dht.network.Chan.Outgoing <- mes
254 + return nil
255 +}
256 +
257 func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *DHTMessage) {
258 dskey := ds.NewKey(pmes.GetKey())
259 var resp *pDHTMessage
@@ -351,10 +384,10 @@ func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *DHTMessage) {
384
385 // Register a handler for a specific message ID, used for getting replies
386 // to certain messages (i.e. response to a GET_VALUE message)
354 -func (dht *IpfsDHT) ListenFor(mesid uint64, count int) <-chan *swarm.Message {
387 +func (dht *IpfsDHT) ListenFor(mesid uint64, count int, timeout time.Duration) <-chan *swarm.Message {
388 lchan := make(chan *swarm.Message)
389 dht.listenLock.Lock()
357 - dht.listeners[mesid] = &listenInfo{lchan, count}
390 + dht.listeners[mesid] = &listenInfo{lchan, count, time.Now().Add(timeout)}
391 dht.listenLock.Unlock()
392 return lchan
393 }
@@ -372,8 +405,14 @@ func (dht *IpfsDHT) Unlisten(mesid uint64) {
405
406 func (dht *IpfsDHT) IsListening(mesid uint64) bool {
407 dht.listenLock.RLock()
375 - _,ok := dht.listeners[mesid]
408 + li,ok := dht.listeners[mesid]
409 dht.listenLock.RUnlock()
410 + if time.Now().After(li.eol) {
411 + dht.listenLock.Lock()
412 + delete(dht.listeners, mesid)
413 + dht.listenLock.Unlock()
414 + return false
415 + }
416 return ok
417 }
418
@@ -401,7 +440,7 @@ func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *DHTMessage) {
440 dht.diaglock.Unlock()
441
442 seq := dht.routes[0].NearestPeers(convertPeerID(dht.self.ID), 10)
404 - listen_chan := dht.ListenFor(pmes.GetId(), len(seq))
443 + listen_chan := dht.ListenFor(pmes.GetId(), len(seq), time.Second * 30)
444
445 for _,ps := range seq {
446 mes := swarm.NewMessage(ps, pmes)
@@ -411,6 +450,9 @@ func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *DHTMessage) {
450
451
452 buf := new(bytes.Buffer)
453 + di := dht.getDiagInfo()
454 + buf.Write(di.Marshal())
455 +
456 // NOTE: this shouldnt be a hardcoded value
457 after := time.After(time.Second * 20)
458 count := len(seq)
@@ -420,14 +462,18 @@ func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *DHTMessage) {
462 //Timeout, return what we have
463 goto out
464 case req_resp := <-listen_chan:
465 + pmes_out := new(DHTMessage)
466 + err := proto.Unmarshal(req_resp.Data, pmes_out)
467 + if err != nil {
468 + // It broke? eh, whatever, keep going
469 + continue
470 + }
471 buf.Write(req_resp.Data)
472 count--
473 }
474 }
475
476 out:
429 - di := dht.getDiagInfo()
430 - buf.Write(di.Marshal())
477 resp := pDHTMessage{
478 Type: DHTMessage_DIAGNOSTIC,
479 Id: pmes.GetId(),
routing/dht/diag.go new
+44
@@ -0,0 +1,44 @@
1 +package dht
2 +
3 +import (
4 + "encoding/json"
5 + "time"
6 +
7 + peer "github.com/jbenet/go-ipfs/peer"
8 +)
9 +
10 +type connDiagInfo struct {
11 + Latency time.Duration
12 + Id peer.ID
13 +}
14 +
15 +type diagInfo struct {
16 + Id peer.ID
17 + Connections []connDiagInfo
18 + Keys []string
19 + LifeSpan time.Duration
20 + CodeVersion string
21 +}
22 +
23 +func (di *diagInfo) Marshal() []byte {
24 + b, err := json.Marshal(di)
25 + if err != nil {
26 + panic(err)
27 + }
28 + //TODO: also consider compressing this. There will be a lot of these
29 + return b
30 +}
31 +
32 +
33 +func (dht *IpfsDHT) getDiagInfo() *diagInfo {
34 + di := new(diagInfo)
35 + di.CodeVersion = "github.com/jbenet/go-ipfs"
36 + di.Id = dht.self.ID
37 + di.LifeSpan = time.Since(dht.birth)
38 + di.Keys = nil // Currently no way to query datastore
39 +
40 + for _,p := range dht.routes[0].listpeers() {
41 + di.Connections = append(di.Connections, connDiagInfo{p.GetDistance(), p.ID})
42 + }
43 + return di
44 +}
routing/dht/routing.go
+14 -16
@@ -29,6 +29,7 @@ func GenerateMessageID() uint64 {
29 // Basic Put/Get
30
31 // PutValue adds value corresponding to given Key.
32 +// This is the top level "Store" operation of the DHT
33 func (s *IpfsDHT) PutValue(key u.Key, value []byte) error {
34 var p *peer.Peer
35 p = s.routes[0].NearestPeer(convertKey(key))
@@ -36,16 +37,7 @@ func (s *IpfsDHT) PutValue(key u.Key, value []byte) error {
37 panic("Table returned nil peer!")
38 }
39
39 - pmes := pDHTMessage{
40 - Type: DHTMessage_PUT_VALUE,
41 - Key: string(key),
42 - Value: value,
43 - Id: GenerateMessageID(),
44 - }
45 -
46 - mes := swarm.NewMessage(p, pmes.ToProtobuf())
47 - s.network.Chan.Outgoing <- mes
48 - return nil
40 + return s.putValueToPeer(p, string(key), value)
41 }
42
43 // GetValue searches for the value corresponding to given Key.
@@ -63,7 +55,7 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
55 Key: string(key),
56 Id: GenerateMessageID(),
57 }
66 - response_chan := s.ListenFor(pmes.Id, 1)
58 + response_chan := s.ListenFor(pmes.Id, 1, time.Minute)
59
60 mes := swarm.NewMessage(p, pmes.ToProtobuf())
61 s.network.Chan.Outgoing <- mes
@@ -74,7 +66,13 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
66 case <-timeup:
67 s.Unlisten(pmes.Id)
68 return nil, u.ErrTimeout
77 - case resp := <-response_chan:
69 + case resp, ok := <-response_chan:
70 + if !ok {
71 + panic("Channel was closed...")
72 + }
73 + if resp == nil {
74 + panic("Why the hell is this response nil?")
75 + }
76 pmes_out := new(DHTMessage)
77 err := proto.Unmarshal(resp.Data, pmes_out)
78 if err != nil {
@@ -123,7 +121,7 @@ func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer,
121
122 mes := swarm.NewMessage(p, pmes.ToProtobuf())
123
126 - listen_chan := s.ListenFor(pmes.Id, 1)
124 + listen_chan := s.ListenFor(pmes.Id, 1, time.Minute)
125 u.DOut("Find providers for: '%s'", key)
126 s.network.Chan.Outgoing <-mes
127 after := time.After(timeout)
@@ -181,7 +179,7 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
179
180 mes := swarm.NewMessage(p, pmes.ToProtobuf())
181
184 - listen_chan := s.ListenFor(pmes.Id, 1)
182 + listen_chan := s.ListenFor(pmes.Id, 1, time.Minute)
183 s.network.Chan.Outgoing <-mes
184 after := time.After(timeout)
185 select {
@@ -224,7 +222,7 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
222 mes := swarm.NewMessage(p, pmes.ToProtobuf())
223
224 before := time.Now()
227 - response_chan := dht.ListenFor(pmes.Id, 1)
225 + response_chan := dht.ListenFor(pmes.Id, 1, time.Minute)
226 dht.network.Chan.Outgoing <- mes
227
228 tout := time.After(timeout)
@@ -253,7 +251,7 @@ func (dht *IpfsDHT) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
251 Id: GenerateMessageID(),
252 }
253
256 - listen_chan := dht.ListenFor(pmes.Id, len(targets))
254 + listen_chan := dht.ListenFor(pmes.Id, len(targets), time.Minute * 2)
255
256 pbmes := pmes.ToProtobuf()
257 for _,p := range targets {
routing/dht/table_test.go
+17
@@ -107,3 +107,20 @@ func TestTableFind(t *testing.T) {
107 t.Fatalf("Failed to lookup known node...")
108 }
109 }
110 +
111 +func TestTableFindMultiple(t *testing.T) {
112 + local := _randPeer()
113 + rt := NewRoutingTable(20, convertPeerID(local.ID))
114 +
115 + peers := make([]*peer.Peer, 100)
116 + for i := 0; i < 18; i++ {
117 + peers[i] = _randPeer()
118 + rt.Update(peers[i])
119 + }
120 +
121 + t.Logf("Searching for peer: '%s'", peers[2].ID.Pretty())
122 + found := rt.NearestPeers(convertPeerID(peers[2].ID), 15)
123 + if len(found) != 15 {
124 + t.Fatalf("Got back different number of peers than we expected.")
125 + }
126 +}