fix bug in diagnostics, and add more peers to closer peer responses
Jeromy committed
Oct 12, 2014 at 23:07 UTC
4189d50d7799fab1bafda077dfe123ad176b3378
4 files changed
+103
-69
core/commands/diag.go
+15
-9
@@ -8,8 +8,22 @@ import (
8
"time"
9
10
"github.com/jbenet/go-ipfs/core"
11
+ diagn "github.com/jbenet/go-ipfs/diagnostics"
12
)
13
14
+func PrintDiagnostics(info []*diagn.DiagInfo, out io.Writer) {
15
+ for _, i := range info {
16
+ fmt.Fprintf(out, "Peer: %s\n", i.ID)
17
+ fmt.Fprintf(out, "\tUp for: %s\n", i.LifeSpan.String())
18
+ fmt.Fprintf(out, "\tConnected To:\n")
19
+ for _, c := range i.Connections {
20
+ fmt.Fprintf(out, "\t%s\n\t\tLatency = %s\n", c.ID, c.Latency.String())
21
+ }
22
+ fmt.Fprintln(out)
23
+ }
24
+
25
+}
26
+
27
func Diag(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.Writer) error {
28
if n.Diagnostics == nil {
29
return errors.New("Cannot run diagnostic in offline mode!")
@@ -29,15 +43,7 @@ func Diag(n *core.IpfsNode, args []string, opts map[string]interface{}, out io.W
43
return err
44
}
45
} else {
32
- for _, i := range info {
33
- fmt.Fprintf(out, "Peer: %s\n", i.ID)
34
- fmt.Fprintf(out, "\tUp for: %s\n", i.LifeSpan.String())
35
- fmt.Fprintf(out, "\tConnected To:\n")
36
- for _, c := range i.Connections {
37
- fmt.Fprintf(out, "\t%s\n\t\tLatency = %s\n", c.ID, c.Latency.String())
38
- }
39
- fmt.Fprintln(out)
40
- }
46
+ PrintDiagnostics(info, out)
47
}
48
return nil
49
}
diagnostics/diag.go
+10
-9
@@ -48,7 +48,7 @@ type connDiagInfo struct {
48
ID string
49
}
50
51
-type diagInfo struct {
51
+type DiagInfo struct {
52
ID string
53
Connections []connDiagInfo
54
Keys []string
@@ -56,7 +56,7 @@ type diagInfo struct {
56
CodeVersion string
57
}
58
59
-func (di *diagInfo) Marshal() []byte {
59
+func (di *DiagInfo) Marshal() []byte {
60
b, err := json.Marshal(di)
61
if err != nil {
62
panic(err)
@@ -69,8 +69,8 @@ func (d *Diagnostics) getPeers() []*peer.Peer {
69
return d.network.GetPeerList()
70
}
71
72
-func (d *Diagnostics) getDiagInfo() *diagInfo {
73
- di := new(diagInfo)
72
+func (d *Diagnostics) getDiagInfo() *DiagInfo {
73
+ di := new(DiagInfo)
74
di.CodeVersion = "github.com/jbenet/go-ipfs"
75
di.ID = d.self.ID.Pretty()
76
di.LifeSpan = time.Since(d.birth)
@@ -88,7 +88,7 @@ func newID() string {
88
return string(id)
89
}
90
91
-func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
91
+func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*DiagInfo, error) {
92
log.Debug("Getting diagnostic.")
93
ctx, _ := context.WithTimeout(context.TODO(), timeout)
94
@@ -102,7 +102,7 @@ func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error)
102
peers := d.getPeers()
103
log.Debug("Sending diagnostic request to %d peers.", len(peers))
104
105
- var out []*diagInfo
105
+ var out []*DiagInfo
106
di := d.getDiagInfo()
107
out = append(out, di)
108
@@ -134,15 +134,15 @@ func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*diagInfo, error)
134
return out, nil
135
}
136
137
-func AppendDiagnostics(data []byte, cur []*diagInfo) []*diagInfo {
137
+func AppendDiagnostics(data []byte, cur []*DiagInfo) []*DiagInfo {
138
buf := bytes.NewBuffer(data)
139
dec := json.NewDecoder(buf)
140
for {
141
- di := new(diagInfo)
141
+ di := new(DiagInfo)
142
err := dec.Decode(di)
143
if err != nil {
144
if err != io.EOF {
145
- log.Error("error decoding diagInfo: %v", err)
145
+ log.Error("error decoding DiagInfo: %v", err)
146
}
147
break
148
}
@@ -216,6 +216,7 @@ func (d *Diagnostics) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, e
216
sendcount := 0
217
for _, p := range d.getPeers() {
218
log.Debug("Sending diagnostic request to peer: %s", p)
219
+ sendcount++
220
go func(p *peer.Peer) {
221
out, err := d.getDiagnosticFromPeer(ctx, p, pmes)
222
if err != nil {
routing/dht/dht.go
+56
-37
@@ -76,7 +76,7 @@ func NewDHT(p *peer.Peer, ps peer.Peerstore, net inet.Network, sender inet.Sende
76
77
// Connect to a new peer at the given address, ping and add to the routing table
78
func (dht *IpfsDHT) Connect(ctx context.Context, npeer *peer.Peer) (*peer.Peer, error) {
79
- log.Debug("Connect to new peer: %s\n", npeer)
79
+ log.Debug("Connect to new peer: %s", npeer)
80
81
// TODO(jbenet,whyrusleeping)
82
//
@@ -109,13 +109,13 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
109
110
mData := mes.Data()
111
if mData == nil {
112
- // TODO handle/log err
112
+ log.Error("Message contained nil data.")
113
return nil
114
}
115
116
mPeer := mes.Peer()
117
if mPeer == nil {
118
- // TODO handle/log err
118
+ log.Error("Message contained nil peer.")
119
return nil
120
}
121
@@ -123,7 +123,7 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
123
pmes := new(Message)
124
err := proto.Unmarshal(mData, pmes)
125
if err != nil {
126
- // TODO handle/log err
126
+ log.Error("Error unmarshaling data")
127
return nil
128
}
129
@@ -138,25 +138,27 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) msg.N
138
handler := dht.handlerForMsgType(pmes.GetType())
139
if handler == nil {
140
// TODO handle/log err
141
+ log.Error("got back nil handler from handlerForMsgType")
142
return nil
143
}
144
145
// dispatch handler.
146
rpmes, err := handler(mPeer, pmes)
147
if err != nil {
147
- // TODO handle/log err
148
+ log.Error("handle message error: %s", err)
149
return nil
150
}
151
152
// if nil response, return it before serializing
153
if rpmes == nil {
154
+ log.Warning("Got back nil response from request.")
155
return nil
156
}
157
158
// serialize response msg
159
rmes, err := msg.FromObject(mPeer, rpmes)
160
if err != nil {
159
- // TODO handle/log err
161
+ log.Error("serialze response error: %s", err)
162
return nil
163
}
164
@@ -197,6 +199,7 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message
199
return rpmes, nil
200
}
201
202
+// putValueToNetwork stores the given key/value pair at the peer 'p'
203
func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p *peer.Peer,
204
key string, value []byte) error {
205
@@ -226,7 +229,7 @@ func (dht *IpfsDHT) putProvider(ctx context.Context, p *peer.Peer, key string) e
229
}
230
231
log.Debug("%s putProvider: %s for %s", dht.self, p, key)
229
- if *rpmes.Key != *pmes.Key {
232
+ if rpmes.GetKey() != pmes.GetKey() {
233
return errors.New("provider not added correctly")
234
}
235
@@ -261,23 +264,11 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
264
// Perhaps we were given closer peers
265
var peers []*peer.Peer
266
for _, pb := range pmes.GetCloserPeers() {
264
- if peer.ID(pb.GetId()).Equal(dht.self.ID) {
265
- continue
266
- }
267
-
268
- addr, err := ma.NewMultiaddr(pb.GetAddr())
267
+ pr, err := dht.addPeer(pb)
268
if err != nil {
270
- log.Error("%v", err.Error())
269
+ log.Error("%s", err)
270
continue
271
}
273
-
274
- // check if we already have this peer.
275
- pr, _ := dht.peerstore.Get(peer.ID(pb.GetId()))
276
- if pr == nil {
277
- pr = &peer.Peer{ID: peer.ID(pb.GetId())}
278
- dht.peerstore.Put(pr)
279
- }
280
- pr.AddAddress(addr) // idempotent
272
peers = append(peers, pr)
273
}
274
@@ -290,6 +281,27 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
281
return nil, nil, u.ErrNotFound
282
}
283
284
+func (dht *IpfsDHT) addPeer(pb *Message_Peer) (*peer.Peer, error) {
285
+ if peer.ID(pb.GetId()).Equal(dht.self.ID) {
286
+ return nil, errors.New("cannot add self as peer")
287
+ }
288
+
289
+ addr, err := ma.NewMultiaddr(pb.GetAddr())
290
+ if err != nil {
291
+ return nil, err
292
+ }
293
+
294
+ // check if we already have this peer.
295
+ pr, _ := dht.peerstore.Get(peer.ID(pb.GetId()))
296
+ if pr == nil {
297
+ pr = &peer.Peer{ID: peer.ID(pb.GetId())}
298
+ dht.peerstore.Put(pr)
299
+ }
300
+ pr.AddAddress(addr) // idempotent
301
+
302
+ return pr, nil
303
+}
304
+
305
// getValueSingle simply performs the get value RPC with the given parameters
306
func (dht *IpfsDHT) getValueSingle(ctx context.Context, p *peer.Peer,
307
key u.Key, level int) (*Message, error) {
@@ -327,6 +339,7 @@ func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
339
return nil, u.ErrNotFound
340
}
341
342
+// getLocal attempts to retrieve the value from the datastore
343
func (dht *IpfsDHT) getLocal(key u.Key) ([]byte, error) {
344
dht.dslock.Lock()
345
defer dht.dslock.Unlock()
@@ -342,6 +355,7 @@ func (dht *IpfsDHT) getLocal(key u.Key) ([]byte, error) {
355
return byt, nil
356
}
357
358
+// putLocal stores the key value pair in the datastore
359
func (dht *IpfsDHT) putLocal(key u.Key, value []byte) error {
360
return dht.datastore.Put(key.DsKey(), value)
361
}
@@ -419,39 +433,44 @@ func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []*peer.Peer
433
return provArr
434
}
435
422
-// nearestPeerToQuery returns the routing tables closest peers.
423
-func (dht *IpfsDHT) nearestPeerToQuery(pmes *Message) *peer.Peer {
436
+// nearestPeersToQuery returns the routing tables closest peers.
437
+func (dht *IpfsDHT) nearestPeersToQuery(pmes *Message, count int) []*peer.Peer {
438
level := pmes.GetClusterLevel()
439
cluster := dht.routingTables[level]
440
441
key := u.Key(pmes.GetKey())
428
- closer := cluster.NearestPeer(kb.ConvertKey(key))
442
+ closer := cluster.NearestPeers(kb.ConvertKey(key), count)
443
return closer
444
}
445
432
-// betterPeerToQuery returns nearestPeerToQuery, but iff closer than self.
433
-func (dht *IpfsDHT) betterPeerToQuery(pmes *Message) *peer.Peer {
434
- closer := dht.nearestPeerToQuery(pmes)
446
+// betterPeerToQuery returns nearestPeersToQuery, but iff closer than self.
447
+func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []*peer.Peer {
448
+ closer := dht.nearestPeersToQuery(pmes, count)
449
450
// no node? nil
451
if closer == nil {
452
return nil
453
}
454
441
- // == to self? nil
442
- if closer.ID.Equal(dht.self.ID) {
443
- log.Error("Attempted to return self! this shouldnt happen...")
444
- return nil
455
+ // == to self? thats bad
456
+ for _, p := range closer {
457
+ if p.ID.Equal(dht.self.ID) {
458
+ log.Error("Attempted to return self! this shouldnt happen...")
459
+ return nil
460
+ }
461
}
462
447
- // self is closer? nil
448
- key := u.Key(pmes.GetKey())
449
- if kb.Closer(dht.self.ID, closer.ID, key) {
450
- return nil
463
+ var filtered []*peer.Peer
464
+ for _, p := range closer {
465
+ // must all be closer than self
466
+ key := u.Key(pmes.GetKey())
467
+ if !kb.Closer(dht.self.ID, p.ID, key) {
468
+ filtered = append(filtered, p)
469
+ }
470
}
471
453
- // ok seems like a closer node.
454
- return closer
472
+ // ok seems like closer nodes
473
+ return filtered
474
}
475
476
func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
routing/dht/handlers.go
+22
-14
@@ -13,6 +13,8 @@ import (
13
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
14
)
15
16
+var CloserPeerCount = 4
17
+
18
// dhthandler specifies the signature of functions that handle DHT messages.
19
type dhtHandler func(*peer.Peer, *Message) (*Message, error)
20
@@ -83,10 +85,12 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
85
}
86
87
// Find closest peer on given cluster to desired key and reply with that info
86
- closer := dht.betterPeerToQuery(pmes)
88
+ closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
89
if closer != nil {
88
- log.Debug("handleGetValue returning a closer peer: '%s'\n", closer)
89
- resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
90
+ for _, p := range closer {
91
+ log.Debug("handleGetValue returning closer peer: '%s'", p)
92
+ }
93
+ resp.CloserPeers = peersToPBPeers(closer)
94
}
95
96
return resp, nil
@@ -109,27 +113,31 @@ func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
113
114
func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error) {
115
resp := newMessage(pmes.GetType(), "", pmes.GetClusterLevel())
112
- var closest *peer.Peer
116
+ var closest []*peer.Peer
117
118
// if looking for self... special case where we send it on CloserPeers.
119
if peer.ID(pmes.GetKey()).Equal(dht.self.ID) {
116
- closest = dht.self
120
+ closest = []*peer.Peer{dht.self}
121
} else {
118
- closest = dht.betterPeerToQuery(pmes)
122
+ closest = dht.betterPeersToQuery(pmes, CloserPeerCount)
123
}
124
125
if closest == nil {
122
- log.Error("handleFindPeer: could not find anything.\n")
126
+ log.Error("handleFindPeer: could not find anything.")
127
return resp, nil
128
}
129
126
- if len(closest.Addresses) == 0 {
127
- log.Error("handleFindPeer: no addresses for connected peer...\n")
128
- return resp, nil
130
+ var withAddresses []*peer.Peer
131
+ for _, p := range closest {
132
+ if len(p.Addresses) > 0 {
133
+ withAddresses = append(withAddresses, p)
134
+ }
135
}
136
131
- log.Debug("handleFindPeer: sending back '%s'\n", closest)
132
- resp.CloserPeers = peersToPBPeers([]*peer.Peer{closest})
137
+ for _, p := range withAddresses {
138
+ log.Debug("handleFindPeer: sending back '%s'", p)
139
+ }
140
+ resp.CloserPeers = peersToPBPeers(withAddresses)
141
return resp, nil
142
}
143
@@ -157,9 +165,9 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) (*Message, e
165
}
166
167
// Also send closer peers.
160
- closer := dht.betterPeerToQuery(pmes)
168
+ closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
169
if closer != nil {
162
- resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
170
+ resp.CloserPeers = peersToPBPeers(closer)
171
}
172
173
return resp, nil