fix a few infinitely looping RPCs
Jeromy committed
Aug 14, 2014 at 08:32 UTC
60d061cb4942b630ab4b12bc74b03f6366df4b78
4 files changed
+165
-51
routing/dht/dht.go
+79
-3
@@ -244,12 +244,14 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
244
}
245
iVal, err := dht.datastore.Get(dskey)
246
if err == nil {
247
+ u.DOut("handleGetValue success!")
248
resp.Success = true
249
resp.Value = iVal.([]byte)
250
} else if err == ds.ErrNotFound {
251
// Check if we know any providers for the requested value
252
provs, ok := dht.providers[u.Key(pmes.GetKey())]
253
if ok && len(provs) > 0 {
254
+ u.DOut("handleGetValue returning %d provider[s]", len(provs))
255
for _, prov := range provs {
256
resp.Peers = append(resp.Peers, prov.Value)
257
}
@@ -265,13 +267,21 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
267
} else {
268
level = int(pmes.GetValue()[0]) // Using value field to specify cluster level
269
}
270
+ u.DOut("handleGetValue searching level %d clusters", level)
271
272
closer := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
273
274
+ if closer.ID.Equal(dht.self.ID) {
275
+ u.DOut("Attempted to return self! this shouldnt happen...")
276
+ resp.Peers = nil
277
+ goto out
278
+ }
279
// If this peer is closer than the one from the table, return nil
280
if kb.Closer(dht.self.ID, closer.ID, u.Key(pmes.GetKey())) {
281
resp.Peers = nil
282
+ u.DOut("handleGetValue could not find a closer node than myself.")
283
} else {
284
+ u.DOut("handleGetValue returning a closer peer: '%s'", closer.ID.Pretty())
285
resp.Peers = []*peer.Peer{closer}
286
}
287
}
@@ -280,6 +290,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
290
panic(err)
291
}
292
293
+out:
294
mes := swarm.NewMessage(p, resp.ToProtobuf())
295
dht.network.Send(mes)
296
}
@@ -349,9 +360,17 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
360
providers := dht.providers[u.Key(pmes.GetKey())]
361
dht.providerLock.RUnlock()
362
if providers == nil || len(providers) == 0 {
352
- // TODO: work on tiering this
353
- closer := dht.routes[0].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
354
- resp.Peers = []*peer.Peer{closer}
363
+ level := 0
364
+ if len(pmes.GetValue()) > 0 {
365
+ level = int(pmes.GetValue()[0])
366
+ }
367
+
368
+ closer := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
369
+ if kb.Closer(dht.self.ID, closer.ID, u.Key(pmes.GetKey())) {
370
+ resp.Peers = nil
371
+ } else {
372
+ resp.Peers = []*peer.Peer{closer}
373
+ }
374
} else {
375
for _, prov := range providers {
376
resp.Peers = append(resp.Peers, prov.Value)
@@ -626,3 +645,60 @@ func (dht *IpfsDHT) PrintTables() {
645
route.Print()
646
}
647
}
648
+
649
+func (dht *IpfsDHT) findProvidersSingle(p *peer.Peer, key u.Key, level int, timeout time.Duration) (*PBDHTMessage, error) {
650
+ pmes := DHTMessage{
651
+ Type: PBDHTMessage_GET_PROVIDERS,
652
+ Key: string(key),
653
+ Id: GenerateMessageID(),
654
+ Value: []byte{byte(level)},
655
+ }
656
+
657
+ mes := swarm.NewMessage(p, pmes.ToProtobuf())
658
+
659
+ listenChan := dht.ListenFor(pmes.Id, 1, time.Minute)
660
+ dht.network.Send(mes)
661
+ after := time.After(timeout)
662
+ select {
663
+ case <-after:
664
+ dht.Unlisten(pmes.Id)
665
+ return nil, u.ErrTimeout
666
+ case resp := <-listenChan:
667
+ u.DOut("FindProviders: got response.")
668
+ pmes_out := new(PBDHTMessage)
669
+ err := proto.Unmarshal(resp.Data, pmes_out)
670
+ if err != nil {
671
+ return nil, err
672
+ }
673
+
674
+ return pmes_out, nil
675
+ }
676
+}
677
+
678
+func (dht *IpfsDHT) addPeerList(key u.Key, peers []*PBDHTMessage_PBPeer) []*peer.Peer {
679
+ var prov_arr []*peer.Peer
680
+ for _, prov := range peers {
681
+ // Dont add outselves to the list
682
+ if peer.ID(prov.GetId()).Equal(dht.self.ID) {
683
+ continue
684
+ }
685
+ // Dont add someone who is already on the list
686
+ p := dht.network.Find(u.Key(prov.GetId()))
687
+ if p == nil {
688
+ u.DOut("given provider %s was not in our network already.", peer.ID(prov.GetId()).Pretty())
689
+ maddr, err := ma.NewMultiaddr(prov.GetAddr())
690
+ if err != nil {
691
+ u.PErr("error connecting to new peer: %s", err)
692
+ continue
693
+ }
694
+ p, err = dht.network.GetConnection(peer.ID(prov.GetId()), maddr)
695
+ if err != nil {
696
+ u.PErr("error connecting to new peer: %s", err)
697
+ continue
698
+ }
699
+ }
700
+ dht.addProviderEntry(key, p)
701
+ prov_arr = append(prov_arr, p)
702
+ }
703
+ return prov_arr
704
+}
routing/dht/dht_logger.go
new
+38
@@ -0,0 +1,38 @@
1
+package dht
2
+
3
+import (
4
+ "encoding/json"
5
+ "time"
6
+
7
+ u "github.com/jbenet/go-ipfs/util"
8
+)
9
+
10
+type logDhtRpc struct {
11
+ Type string
12
+ Start time.Time
13
+ End time.Time
14
+ Duration time.Duration
15
+ RpcCount int
16
+ Success bool
17
+}
18
+
19
+func startNewRpc(name string) *logDhtRpc {
20
+ r := new(logDhtRpc)
21
+ r.Type = name
22
+ r.Start = time.Now()
23
+ return r
24
+}
25
+
26
+func (l *logDhtRpc) EndLog() {
27
+ l.End = time.Now()
28
+ l.Duration = l.End.Sub(l.Start)
29
+}
30
+
31
+func (l *logDhtRpc) Print() {
32
+ b, err := json.Marshal(l)
33
+ if err != nil {
34
+ u.POut(err.Error())
35
+ } else {
36
+ u.POut(string(b))
37
+ }
38
+}
routing/dht/dht_test.go
+1
-1
@@ -156,7 +156,7 @@ func TestValueGetSet(t *testing.T) {
156
}
157
158
if string(val) != "world" {
159
- t.Fatalf("Expected 'world' got %s", string(val))
159
+ t.Fatalf("Expected 'world' got '%s'", string(val))
160
}
161
}
162
routing/dht/routing.go
+47
-47
@@ -60,12 +60,19 @@ func (s *IpfsDHT) PutValue(key u.Key, value []byte) {
60
// If the search does not succeed, a multiaddr string of a closer peer is
61
// returned along with util.ErrSearchIncomplete
62
func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
63
+ ll := startNewRpc("GET")
64
+ defer func() {
65
+ ll.EndLog()
66
+ ll.Print()
67
+ }()
68
route_level := 0
69
70
// If we have it local, dont bother doing an RPC!
71
// NOTE: this might not be what we want to do...
67
- val,err := s.GetLocal(key)
68
- if err != nil {
72
+ val, err := s.GetLocal(key)
73
+ if err == nil {
74
+ ll.Success = true
75
+ u.DOut("Found local, returning.")
76
return val, nil
77
}
78
@@ -74,11 +81,8 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
81
return nil, kb.ErrLookupFailure
82
}
83
77
- if kb.Closer(s.self.ID, p.ID, key) {
78
- return nil, u.ErrNotFound
79
- }
80
-
84
for route_level < len(s.routes) && p != nil {
85
+ ll.RpcCount++
86
pmes, err := s.getValueSingle(p, key, timeout, route_level)
87
if err != nil {
88
return nil, u.WrapError(err, "getValue Error")
@@ -86,16 +90,19 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
90
91
if pmes.GetSuccess() {
92
if pmes.Value == nil { // We were given provider[s]
93
+ ll.RpcCount++
94
return s.getFromPeerList(key, timeout, pmes.GetPeers(), route_level)
95
}
96
97
// Success! We were given the value
98
+ ll.Success = true
99
return pmes.GetValue(), nil
100
} else {
101
// We were given a closer node
102
closers := pmes.GetPeers()
103
if len(closers) > 0 {
104
if peer.ID(closers[0].GetId()).Equal(s.self.ID) {
105
+ u.DOut("Got myself back as a closer peer.")
106
return nil, u.ErrNotFound
107
}
108
maddr, err := ma.NewMultiaddr(closers[0].GetAddr())
@@ -108,6 +115,7 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
115
if err != nil {
116
u.PErr("[%s] Failed to connect to: %s", s.self.ID.Pretty(), closers[0].GetAddr())
117
route_level++
118
+ continue
119
}
120
p = np
121
} else {
@@ -143,60 +151,52 @@ func (s *IpfsDHT) Provide(key u.Key) error {
151
152
// FindProviders searches for peers who can provide the value for given key.
153
func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
154
+ ll := startNewRpc("FindProviders")
155
+ defer func() {
156
+ ll.EndLog()
157
+ ll.Print()
158
+ }()
159
+ u.DOut("Find providers for: '%s'", key)
160
p := s.routes[0].NearestPeer(kb.ConvertKey(key))
161
if p == nil {
162
return nil, kb.ErrLookupFailure
163
}
164
151
- pmes := DHTMessage{
152
- Type: PBDHTMessage_GET_PROVIDERS,
153
- Key: string(key),
154
- Id: GenerateMessageID(),
155
- }
156
-
157
- mes := swarm.NewMessage(p, pmes.ToProtobuf())
158
-
159
- listenChan := s.ListenFor(pmes.Id, 1, time.Minute)
160
- u.DOut("Find providers for: '%s'", key)
161
- s.network.Send(mes)
162
- after := time.After(timeout)
163
- select {
164
- case <-after:
165
- s.Unlisten(pmes.Id)
166
- return nil, u.ErrTimeout
167
- case resp := <-listenChan:
168
- u.DOut("FindProviders: got response.")
169
- pmes_out := new(PBDHTMessage)
170
- err := proto.Unmarshal(resp.Data, pmes_out)
165
+ for level := 0; level < len(s.routes); {
166
+ pmes, err := s.findProvidersSingle(p, key, level, timeout)
167
if err != nil {
168
return nil, err
169
}
174
-
175
- var prov_arr []*peer.Peer
176
- for _, prov := range pmes_out.GetPeers() {
177
- if peer.ID(prov.GetId()).Equal(s.self.ID) {
170
+ if pmes.GetSuccess() {
171
+ provs := s.addPeerList(key, pmes.GetPeers())
172
+ ll.Success = true
173
+ return provs, nil
174
+ } else {
175
+ closer := pmes.GetPeers()
176
+ if len(closer) == 0 {
177
+ level++
178
continue
179
}
180
- p := s.network.Find(u.Key(prov.GetId()))
181
- if p == nil {
182
- u.DOut("given provider %s was not in our network already.", peer.ID(prov.GetId()).Pretty())
183
- maddr, err := ma.NewMultiaddr(prov.GetAddr())
184
- if err != nil {
185
- u.PErr("error connecting to new peer: %s", err)
186
- continue
187
- }
188
- p, err = s.network.GetConnection(peer.ID(prov.GetId()), maddr)
189
- if err != nil {
190
- u.PErr("error connecting to new peer: %s", err)
191
- continue
192
- }
180
+ if peer.ID(closer[0].GetId()).Equal(s.self.ID) {
181
+ u.DOut("Got myself back as a closer peer.")
182
+ return nil, u.ErrNotFound
183
+ }
184
+ maddr, err := ma.NewMultiaddr(closer[0].GetAddr())
185
+ if err != nil {
186
+ // ??? Move up route level???
187
+ panic("not yet implemented")
188
}
194
- s.addProviderEntry(key, p)
195
- prov_arr = append(prov_arr, p)
196
- }
189
198
- return prov_arr, nil
190
+ np, err := s.network.GetConnection(peer.ID(closer[0].GetId()), maddr)
191
+ if err != nil {
192
+ u.PErr("[%s] Failed to connect to: %s", s.self.ID.Pretty(), closer[0].GetAddr())
193
+ level++
194
+ continue
195
+ }
196
+ p = np
197
+ }
198
}
199
+ return nil, u.ErrNotFound
200
}
201
202
// Find specific Peer