@cryptotaxi247 / kubo / commits / 67bd041b9

got everything to build

Juan Batiz-Benet committed Sep 17, 2014 at 07:19 UTC 67bd041b9ce4b1141f90630627c7033c0cb11439
5 files changed +270 -253
routing/dht/dht.go
+27 -1
@@ -142,7 +142,6 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) (msg.
142 Message_MessageType_name[int32(pmes.GetType())], mPeer.ID.Pretty())
143
144 // get handler for this msg type.
145 - var resp *Message
145 handler := dht.handlerForMsgType(pmes.GetType())
146 if handler == nil {
147 return nil, errors.New("Recieved invalid message type")
@@ -190,6 +189,27 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message
189 return rpmes, nil
190 }
191
192 +func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p *peer.Peer, key string, value []byte) error {
193 + pmes := newMessage(Message_PUT_VALUE, string(key), 0)
194 + pmes.Value = value
195 +
196 + mes, err := msg.FromObject(p, pmes)
197 + if err != nil {
198 + return err
199 + }
200 + return dht.sender.SendMessage(ctx, mes)
201 +}
202 +
203 +func (dht *IpfsDHT) putProvider(ctx context.Context, p *peer.Peer, key string) error {
204 + pmes := newMessage(Message_ADD_PROVIDER, string(key), 0)
205 +
206 + mes, err := msg.FromObject(p, pmes)
207 + if err != nil {
208 + return err
209 + }
210 + return dht.sender.SendMessage(ctx, mes)
211 +}
212 +
213 func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
214 key u.Key, level int) ([]byte, []*peer.Peer, error) {
215
@@ -406,6 +426,12 @@ func (dht *IpfsDHT) betterPeerToQuery(pmes *Message) *peer.Peer {
426 func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
427
428 id := peer.ID(pbp.GetId())
429 +
430 + // continue if it's ourselves
431 + if id.Equal(dht.self.ID) {
432 + return nil, errors.New("found self")
433 + }
434 +
435 p, _ := dht.peerstore.Get(id)
436 if p == nil {
437 p, _ = dht.Find(id)
routing/dht/dht_logger.go
+5
@@ -36,3 +36,8 @@ func (l *logDhtRPC) Print() {
36 u.DOut(string(b))
37 }
38 }
39 +
40 +func (l *logDhtRPC) EndAndPrint() {
41 + l.EndLog()
42 + l.Print()
43 +}
routing/dht/handlers.go
+1 -14
@@ -10,7 +10,6 @@ import (
10 kb "github.com/jbenet/go-ipfs/routing/kbucket"
11 u "github.com/jbenet/go-ipfs/util"
12
13 - context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
14 )
15
@@ -38,18 +37,6 @@ func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
37 }
38 }
39
41 -func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
42 - typ := Message_PUT_VALUE
43 - pmes := newMessage(Message_PUT_VALUE, string(key), 0)
44 - pmes.Value = value
45 -
46 - mes, err := msg.FromObject(p, pmes)
47 - if err != nil {
48 - return err
49 - }
50 - return dht.sender.SendMessage(context.TODO(), mes)
51 -}
52 -
40 func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error) {
41 u.DOut("handleGetValue for key: %s\n", pmes.GetKey())
42
@@ -205,7 +192,7 @@ func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *Message) (*Message, err
192 seq := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
193
194 for _, ps := range seq {
208 - mes, err := msg.FromObject(ps, pmes)
195 + _, err := msg.FromObject(ps, pmes)
196 if err != nil {
197 u.PErr("handleDiagnostics error creating message: %v\n", err)
198 continue
routing/dht/query.go new
+85
@@ -0,0 +1,85 @@
1 +package dht
2 +
3 +import (
4 + peer "github.com/jbenet/go-ipfs/peer"
5 + queue "github.com/jbenet/go-ipfs/peer/queue"
6 + u "github.com/jbenet/go-ipfs/util"
7 +
8 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9 +)
10 +
11 +type dhtQuery struct {
12 + // a PeerQueue
13 + peers queue.PeerQueue
14 +
15 + // the function to execute per peer
16 + qfunc queryFunc
17 +}
18 +
19 +// QueryFunc is a function that runs a particular query with a given peer.
20 +// It returns either:
21 +// - the value
22 +// - a list of peers potentially better able to serve the query
23 +// - an error
24 +type queryFunc func(context.Context, *peer.Peer) (interface{}, []*peer.Peer, error)
25 +
26 +func (q *dhtQuery) Run(ctx context.Context, concurrency int) (interface{}, error) {
27 + // get own cancel function to signal when we've found the value
28 + ctx, cancel := context.WithCancel(ctx)
29 +
30 + // the variable waiting to be populated upon success
31 + var result interface{}
32 +
33 + // chanQueue is how workers receive their work
34 + chanQueue := queue.NewChanQueue(ctx, q.peers)
35 +
36 + // worker
37 + worker := func() {
38 + for {
39 + select {
40 + case p := <-chanQueue.DeqChan:
41 +
42 + val, closer, err := q.qfunc(ctx, p)
43 + if err != nil {
44 + u.PErr("error running query: %v\n", err)
45 + continue
46 + }
47 +
48 + if val != nil {
49 + result = val
50 + cancel() // signal we're done.
51 + return
52 + }
53 +
54 + if closer != nil {
55 + for _, p := range closer {
56 + select {
57 + case chanQueue.EnqChan <- p:
58 + case <-ctx.Done():
59 + return
60 + }
61 + }
62 + }
63 +
64 + case <-ctx.Done():
65 + return
66 + }
67 + }
68 + }
69 +
70 + // launch all workers
71 + for i := 0; i < concurrency; i++ {
72 + go worker()
73 + }
74 +
75 + // wait until we're done. yep.
76 + select {
77 + case <-ctx.Done():
78 + }
79 +
80 + if result != nil {
81 + return result, nil
82 + }
83 +
84 + return nil, ctx.Err()
85 +}
routing/dht/routing.go
+152 -238
@@ -4,14 +4,13 @@ import (
4 "bytes"
5 "encoding/json"
6 "errors"
7 + "fmt"
8 "time"
9
9 - proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
10 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11
11 - ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
12 -
13 - swarm "github.com/jbenet/go-ipfs/net/swarm"
12 peer "github.com/jbenet/go-ipfs/peer"
13 + queue "github.com/jbenet/go-ipfs/peer/queue"
14 kb "github.com/jbenet/go-ipfs/routing/kbucket"
15 u "github.com/jbenet/go-ipfs/util"
16 )
@@ -23,29 +22,31 @@ import (
22 // PutValue adds value corresponding to given Key.
23 // This is the top level "Store" operation of the DHT
24 func (dht *IpfsDHT) PutValue(key u.Key, value []byte) error {
26 - complete := make(chan struct{})
27 - count := 0
25 + ctx := context.TODO()
26 +
27 + query := &dhtQuery{}
28 + query.peers = queue.NewXORDistancePQ(key)
29 +
30 + // get the peers we need to announce to
31 for _, route := range dht.routingTables {
32 peers := route.NearestPeers(kb.ConvertKey(key), KValue)
33 for _, p := range peers {
34 if p == nil {
32 - dht.network.Error(kb.ErrLookupFailure)
33 - continue
35 + // this shouldn't be happening.
36 + panic("p should not be nil")
37 }
35 - count++
36 - go func(sp *peer.Peer) {
37 - err := dht.putValueToNetwork(sp, string(key), value)
38 - if err != nil {
39 - dht.network.Error(err)
40 - }
41 - complete <- struct{}{}
42 - }(p)
38 +
39 + query.peers.Enqueue(p)
40 }
41 }
45 - for i := 0; i < count; i++ {
46 - <-complete
42 +
43 + query.qfunc = func(ctx context.Context, p *peer.Peer) (interface{}, []*peer.Peer, error) {
44 + dht.putValueToNetwork(ctx, p, string(key), value)
45 + return nil, nil, nil
46 }
48 - return nil
47 +
48 + _, err := query.Run(ctx, query.peers.Len())
49 + return err
50 }
51
52 // GetValue searches for the value corresponding to given Key.
@@ -53,10 +54,9 @@ func (dht *IpfsDHT) PutValue(key u.Key, value []byte) error {
54 // returned along with util.ErrSearchIncomplete
55 func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
56 ll := startNewRPC("GET")
56 - defer func() {
57 - ll.EndLog()
58 - ll.Print()
59 - }()
57 + defer ll.EndAndPrint()
58 +
59 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
60
61 // If we have it local, dont bother doing an RPC!
62 // NOTE: this might not be what we want to do...
@@ -67,98 +67,37 @@ func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
67 return val, nil
68 }
69
70 + // get closest peers in the routing tables
71 routeLevel := 0
72 closest := dht.routingTables[routeLevel].NearestPeers(kb.ConvertKey(key), PoolSize)
73 if closest == nil || len(closest) == 0 {
74 return nil, kb.ErrLookupFailure
75 }
76
76 - valChan := make(chan []byte)
77 - npeerChan := make(chan *peer.Peer, 30)
78 - procPeer := make(chan *peer.Peer, 30)
79 - errChan := make(chan error)
80 - after := time.After(timeout)
81 - pset := newPeerSet()
77 + query := &dhtQuery{}
78 + query.peers = queue.NewXORDistancePQ(key)
79
80 + // get the peers we need to announce to
81 for _, p := range closest {
84 - pset.Add(p)
85 - npeerChan <- p
82 + query.peers.Enqueue(p)
83 }
84
88 - c := counter{}
89 -
90 - count := 0
91 - go func() {
92 - defer close(procPeer)
93 - for {
94 - select {
95 - case p, ok := <-npeerChan:
96 - if !ok {
97 - return
98 - }
99 - count++
100 - if count >= KValue {
101 - errChan <- u.ErrNotFound
102 - return
103 - }
104 - c.Increment()
105 -
106 - procPeer <- p
107 - default:
108 - if c.Size() <= 0 {
109 - select {
110 - case errChan <- u.ErrNotFound:
111 - default:
112 - }
113 - return
114 - }
115 - }
116 - }
117 - }()
118 -
119 - process := func() {
120 - defer c.Decrement()
121 - for p := range procPeer {
122 - if p == nil {
123 - return
124 - }
125 - val, peers, err := dht.getValueOrPeers(p, key, timeout/4, routeLevel)
126 - if err != nil {
127 - u.DErr("%v\n", err.Error())
128 - continue
129 - }
130 - if val != nil {
131 - select {
132 - case valChan <- val:
133 - default:
134 - u.DOut("Wasnt the first to return the value!")
135 - }
136 - return
137 - }
138 -
139 - for _, np := range peers {
140 - // TODO: filter out peers that arent closer
141 - if !pset.Contains(np) && pset.Size() < KValue {
142 - pset.Add(np) //This is racey... make a single function to do operation
143 - npeerChan <- np
144 - }
145 - }
146 - c.Decrement()
147 - }
85 + // setup the Query Function
86 + query.qfunc = func(ctx context.Context, p *peer.Peer) (interface{}, []*peer.Peer, error) {
87 + return dht.getValueOrPeers(ctx, p, key, routeLevel)
88 }
89
150 - for i := 0; i < AlphaValue; i++ {
151 - go process()
90 + // run it!
91 + result, err := query.Run(ctx, query.peers.Len())
92 + if err != nil {
93 + return nil, err
94 }
95
154 - select {
155 - case val := <-valChan:
156 - return val, nil
157 - case err := <-errChan:
158 - return nil, err
159 - case <-after:
160 - return nil, u.ErrTimeout
96 + byt, ok := result.([]byte)
97 + if !ok {
98 + return nil, fmt.Errorf("received non-byte slice value")
99 }
100 + return byt, nil
101 }
102
103 // Value provider layer of indirection.
@@ -166,26 +105,27 @@ func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
105
106 // Provide makes this node announce that it can provide a value for the given key
107 func (dht *IpfsDHT) Provide(key u.Key) error {
108 + ctx := context.TODO()
109 +
110 dht.providers.AddProvider(key, dht.self)
111 peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), PoolSize)
112 if len(peers) == 0 {
113 return kb.ErrLookupFailure
114 }
115
175 - pmes := Message{
176 - Type: PBDHTMessage_ADD_PROVIDER,
177 - Key: string(key),
178 - }
179 - pbmes := pmes.ToProtobuf()
180 -
116 for _, p := range peers {
182 - mes := swarm.NewMessage(p, pbmes)
183 - dht.netChan.Outgoing <- mes
117 + err := dht.putProvider(ctx, p, string(key))
118 + if err != nil {
119 + return err
120 + }
121 }
122 return nil
123 }
124
125 +// FindProvidersAsync runs FindProviders and sends back results over a channel
126 func (dht *IpfsDHT) FindProvidersAsync(key u.Key, count int, timeout time.Duration) chan *peer.Peer {
127 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
128 +
129 peerOut := make(chan *peer.Peer, count)
130 go func() {
131 ps := newPeerSet()
@@ -202,13 +142,14 @@ func (dht *IpfsDHT) FindProvidersAsync(key u.Key, count int, timeout time.Durati
142
143 peers := dht.routingTables[0].NearestPeers(kb.ConvertKey(key), AlphaValue)
144 for _, pp := range peers {
145 + ppp := pp
146 go func() {
206 - pmes, err := dht.findProvidersSingle(pp, key, 0, timeout)
147 + pmes, err := dht.findProvidersSingle(ctx, ppp, key, 0)
148 if err != nil {
149 u.PErr("%v\n", err)
150 return
151 }
211 - dht.addPeerListAsync(key, pmes.GetPeers(), ps, count, peerOut)
152 + dht.addPeerListAsync(key, pmes.GetProviderPeers(), ps, count, peerOut)
153 }()
154 }
155
@@ -217,21 +158,15 @@ func (dht *IpfsDHT) FindProvidersAsync(key u.Key, count int, timeout time.Durati
158 }
159
160 //TODO: this function could also be done asynchronously
220 -func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*PBDHTMessage_PBPeer, ps *peerSet, count int, out chan *peer.Peer) {
161 +func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*Message_Peer, ps *peerSet, count int, out chan *peer.Peer) {
162 for _, pbp := range peers {
222 - if peer.ID(pbp.GetId()).Equal(dht.self.ID) {
223 - continue
224 - }
225 - maddr, err := ma.NewMultiaddr(pbp.GetAddr())
226 - if err != nil {
227 - u.PErr("%v\n", err)
228 - continue
229 - }
230 - p, err := dht.network.GetConnection(peer.ID(pbp.GetId()), maddr)
163 +
164 + // construct new peer
165 + p, err := dht.ensureConnectedToPeer(pbp)
166 if err != nil {
232 - u.PErr("%v\n", err)
167 continue
168 }
169 +
170 dht.providers.AddProvider(k, p)
171 if ps.AddIfSmallerThan(p, count) {
172 out <- p
@@ -244,10 +179,11 @@ func (dht *IpfsDHT) addPeerListAsync(k u.Key, peers []*PBDHTMessage_PBPeer, ps *
179 // FindProviders searches for peers who can provide the value for given key.
180 func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer, error) {
181 ll := startNewRPC("FindProviders")
247 - defer func() {
248 - ll.EndLog()
249 - ll.Print()
250 - }()
182 + ll.EndAndPrint()
183 +
184 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
185 +
186 + // get closest peer
187 u.DOut("Find providers for: '%s'\n", key)
188 p := dht.routingTables[0].NearestPeer(kb.ConvertKey(key))
189 if p == nil {
@@ -255,37 +191,30 @@ func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Pee
191 }
192
193 for level := 0; level < len(dht.routingTables); {
258 - pmes, err := dht.findProvidersSingle(p, key, level, timeout)
194 +
195 + // attempt retrieving providers
196 + pmes, err := dht.findProvidersSingle(ctx, p, key, level)
197 if err != nil {
198 return nil, err
199 }
262 - if pmes.GetSuccess() {
200 +
201 + // handle providers
202 + provs := pmes.GetProviderPeers()
203 + if provs != nil {
204 u.DOut("Got providers back from findProviders call!\n")
264 - provs := dht.addProviders(key, pmes.GetPeers())
265 - ll.Success = true
266 - return provs, nil
205 + return dht.addProviders(key, provs), nil
206 }
207
208 u.DOut("Didnt get providers, just closer peers.\n")
270 -
271 - closer := pmes.GetPeers()
209 + closer := pmes.GetCloserPeers()
210 if len(closer) == 0 {
211 level++
212 continue
213 }
276 - if peer.ID(closer[0].GetId()).Equal(dht.self.ID) {
277 - u.DOut("Got myself back as a closer peer.")
278 - return nil, u.ErrNotFound
279 - }
280 - maddr, err := ma.NewMultiaddr(closer[0].GetAddr())
281 - if err != nil {
282 - // ??? Move up route level???
283 - panic("not yet implemented")
284 - }
214
286 - np, err := dht.network.GetConnection(peer.ID(closer[0].GetId()), maddr)
215 + np, err := dht.peerFromInfo(closer[0])
216 if err != nil {
288 - u.PErr("[%s] Failed to connect to: %s\n", dht.self.ID.Pretty(), closer[0].GetAddr())
217 + u.DOut("no peerFromInfo")
218 level++
219 continue
220 }
@@ -298,12 +227,15 @@ func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Pee
227
228 // FindPeer searches for a peer with given ID.
229 func (dht *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error) {
230 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
231 +
232 // Check if were already connected to them
233 p, _ := dht.Find(id)
234 if p != nil {
235 return p, nil
236 }
237
238 + // @whyrusleeping why is this here? doesn't the dht.Find above cover it?
239 routeLevel := 0
240 p = dht.routingTables[routeLevel].NearestPeer(kb.ConvertPeerID(id))
241 if p == nil {
@@ -314,158 +246,140 @@ func (dht *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, err
246 }
247
248 for routeLevel < len(dht.routingTables) {
317 - pmes, err := dht.findPeerSingle(p, id, timeout, routeLevel)
318 - plist := pmes.GetPeers()
249 + pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
250 + plist := pmes.GetCloserPeers()
251 if plist == nil || len(plist) == 0 {
252 routeLevel++
253 continue
254 }
255 found := plist[0]
256
325 - addr, err := ma.NewMultiaddr(found.GetAddr())
257 + nxtPeer, err := dht.ensureConnectedToPeer(found)
258 if err != nil {
327 - return nil, err
259 + routeLevel++
260 + continue
261 }
262
330 - nxtPeer, err := dht.network.GetConnection(peer.ID(found.GetId()), addr)
331 - if err != nil {
332 - return nil, err
333 - }
334 - if pmes.GetSuccess() {
335 - if !id.Equal(nxtPeer.ID) {
336 - return nil, errors.New("got back invalid peer from 'successful' response")
337 - }
263 + if nxtPeer.ID.Equal(id) {
264 return nxtPeer, nil
265 }
266 +
267 p = nxtPeer
268 }
269 return nil, u.ErrNotFound
270 }
271
272 func (dht *IpfsDHT) findPeerMultiple(id peer.ID, timeout time.Duration) (*peer.Peer, error) {
273 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
274 +
275 // Check if were already connected to them
276 p, _ := dht.Find(id)
277 if p != nil {
278 return p, nil
279 }
280
281 + query := &dhtQuery{}
282 + query.peers = queue.NewXORDistancePQ(u.Key(id))
283 +
284 + // get the peers we need to announce to
285 routeLevel := 0
286 peers := dht.routingTables[routeLevel].NearestPeers(kb.ConvertPeerID(id), AlphaValue)
287 if len(peers) == 0 {
288 return nil, kb.ErrLookupFailure
289 }
290 + for _, p := range peers {
291 + query.peers.Enqueue(p)
292 + }
293 +
294 + // setup query function
295 + query.qfunc = func(ctx context.Context, p *peer.Peer) (interface{}, []*peer.Peer, error) {
296 + pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
297 + if err != nil {
298 + u.DErr("getPeer error: %v\n", err)
299 + return nil, nil, err
300 + }
301
358 - found := make(chan *peer.Peer)
359 - after := time.After(timeout)
302 + plist := pmes.GetCloserPeers()
303 + if len(plist) == 0 {
304 + routeLevel++
305 + }
306
361 - for _, p := range peers {
362 - go func(p *peer.Peer) {
363 - pmes, err := dht.findPeerSingle(p, id, timeout, routeLevel)
307 + nxtprs := make([]*peer.Peer, len(plist))
308 + for i, fp := range plist {
309 + nxtp, err := dht.peerFromInfo(fp)
310 if err != nil {
365 - u.DErr("getPeer error: %v\n", err)
366 - return
367 - }
368 - plist := pmes.GetPeers()
369 - if len(plist) == 0 {
370 - routeLevel++
311 + u.DErr("findPeer error: %v\n", err)
312 + continue
313 }
372 - for _, fp := range plist {
373 - nxtp, err := dht.peerFromInfo(fp)
374 - if err != nil {
375 - u.DErr("findPeer error: %v\n", err)
376 - continue
377 - }
314
379 - if nxtp.ID.Equal(dht.self.ID) {
380 - found <- nxtp
381 - return
382 - }
315 + if nxtp.ID.Equal(id) {
316 + return nxtp, nil, nil
317 }
384 - }(p)
318 +
319 + nxtprs[i] = nxtp
320 + }
321 +
322 + return nil, nxtprs, nil
323 }
324
387 - select {
388 - case p := <-found:
389 - return p, nil
390 - case <-after:
391 - return nil, u.ErrTimeout
325 + p5, err := query.Run(ctx, query.peers.Len())
326 + if err != nil {
327 + return nil, err
328 + }
329 +
330 + p6, ok := p5.(*peer.Peer)
331 + if !ok {
332 + return nil, errors.New("received non peer object")
333 }
334 + return p6, nil
335 }
336
337 // Ping a peer, log the time it took
338 func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
339 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
340 +
341 // Thoughts: maybe this should accept an ID and do a peer lookup?
342 u.DOut("Enter Ping.\n")
343
400 - pmes := Message{ID: swarm.GenerateMessageID(), Type: PBDHTMessage_PING}
401 - mes := swarm.NewMessage(p, pmes.ToProtobuf())
402 -
403 - before := time.Now()
404 - responseChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
405 - dht.netChan.Outgoing <- mes
406 -
407 - tout := time.After(timeout)
408 - select {
409 - case <-responseChan:
410 - roundtrip := time.Since(before)
411 - p.SetLatency(roundtrip)
412 - u.DOut("Ping took %s.\n", roundtrip.String())
413 - return nil
414 - case <-tout:
415 - // Timed out, think about removing peer from network
416 - u.DOut("[%s] Ping peer [%s] timed out.", dht.self.ID.Pretty(), p.ID.Pretty())
417 - dht.listener.Unlisten(pmes.ID)
418 - return u.ErrTimeout
419 - }
344 + pmes := newMessage(Message_PING, "", 0)
345 + _, err := dht.sendRequest(ctx, p, pmes)
346 + return err
347 }
348
349 func (dht *IpfsDHT) getDiagnostic(timeout time.Duration) ([]*diagInfo, error) {
423 - u.DOut("Begin Diagnostic")
424 - //Send to N closest peers
425 - targets := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
426 -
427 - // TODO: Add timeout to this struct so nodes know when to return
428 - pmes := Message{
429 - Type: PBDHTMessage_DIAGNOSTIC,
430 - ID: swarm.GenerateMessageID(),
431 - }
350 + ctx, _ := context.WithTimeout(context.TODO(), timeout)
351
433 - listenChan := dht.listener.Listen(pmes.ID, len(targets), time.Minute*2)
352 + u.DOut("Begin Diagnostic")
353 + query := &dhtQuery{}
354 + query.peers = queue.NewXORDistancePQ(u.Key(dht.self.ID))
355
435 - pbmes := pmes.ToProtobuf()
356 + targets := dht.routingTables[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
357 for _, p := range targets {
437 - mes := swarm.NewMessage(p, pbmes)
438 - dht.netChan.Outgoing <- mes
358 + query.peers.Enqueue(p)
359 }
360
361 var out []*diagInfo
442 - after := time.After(timeout)
443 - for count := len(targets); count > 0; {
444 - select {
445 - case <-after:
446 - u.DOut("Diagnostic request timed out.")
447 - return out, u.ErrTimeout
448 - case resp := <-listenChan:
449 - pmesOut := new(PBDHTMessage)
450 - err := proto.Unmarshal(resp.Data, pmesOut)
451 - if err != nil {
452 - // NOTE: here and elsewhere, need to audit error handling,
453 - // some errors should be continued on from
454 - return out, err
455 - }
362
457 - dec := json.NewDecoder(bytes.NewBuffer(pmesOut.GetValue()))
458 - for {
459 - di := new(diagInfo)
460 - err := dec.Decode(di)
461 - if err != nil {
462 - break
463 - }
363 + query.qfunc = func(ctx context.Context, p *peer.Peer) (interface{}, []*peer.Peer, error) {
364 + pmes := newMessage(Message_DIAGNOSTIC, "", 0)
365 + rpmes, err := dht.sendRequest(ctx, p, pmes)
366 + if err != nil {
367 + return nil, nil, err
368 + }
369
465 - out = append(out, di)
370 + dec := json.NewDecoder(bytes.NewBuffer(rpmes.GetValue()))
371 + for {
372 + di := new(diagInfo)
373 + err := dec.Decode(di)
374 + if err != nil {
375 + break
376 }
377 +
378 + out = append(out, di)
379 }
380 + return nil, nil, nil
381 }
382
470 - return nil, nil
383 + _, err := query.Run(ctx, query.peers.Len())
384 + return out, err
385 }