query wasnt ensuring conn
The query-- once it's actually attempting to connect to a peer-- should be the one connecting.
Juan Batiz-Benet committed
Oct 21, 2014 at 01:18 UTC
c3df3973e1ff468e9b091f6e948664994835a185
2 files changed
+35
-12
routing/dht/query.go
+32
-9
@@ -14,10 +14,18 @@ import (
14
15
var maxQueryConcurrency = AlphaValue
16
17
+type dhtDialer interface {
18
+ // DialPeer attempts to establish a connection to a given peer
19
+ DialPeer(peer.Peer) error
20
+}
21
+
22
type dhtQuery struct {
23
// the key we're querying for
24
key u.Key
25
26
+ // dialer used to ensure we're connected to peers
27
+ dialer dhtDialer
28
+
29
// the function to execute per peer
30
qfunc queryFunc
31
@@ -34,9 +42,10 @@ type dhtQueryResult struct {
42
}
43
44
// constructs query
37
-func newQuery(k u.Key, f queryFunc) *dhtQuery {
45
+func newQuery(k u.Key, d dhtDialer, f queryFunc) *dhtQuery {
46
return &dhtQuery{
47
key: k,
48
+ dialer: d,
49
qfunc: f,
50
concurrency: maxQueryConcurrency,
51
}
@@ -211,19 +220,38 @@ func (r *dhtQueryRunner) queryPeer(p peer.Peer) {
220
return
221
}
222
214
- log.Debug("running worker for: %v\n", p)
223
+ // ok let's do this!
224
+ log.Debug("running worker for: %v", p)
225
+
226
+ // make sure we do this when we exit
227
+ defer func() {
228
+ // signal we're done proccessing peer p
229
+ log.Debug("completing worker for: %v", p)
230
+ r.peersRemaining.Decrement(1)
231
+ r.rateLimit <- struct{}{}
232
+ }()
233
+
234
+ // make sure we're connected to the peer.
235
+ err := r.query.dialer.DialPeer(p)
236
+ if err != nil {
237
+ log.Debug("ERROR worker for: %v -- err connecting: %v", p, err)
238
+ r.Lock()
239
+ r.errs = append(r.errs, err)
240
+ r.Unlock()
241
+ return
242
+ }
243
244
// finally, run the query against this peer
245
res, err := r.query.qfunc(r.ctx, p)
246
247
if err != nil {
220
- log.Debug("ERROR worker for: %v %v\n", p, err)
248
+ log.Debug("ERROR worker for: %v %v", p, err)
249
r.Lock()
250
r.errs = append(r.errs, err)
251
r.Unlock()
252
253
} else if res.success {
226
- log.Debug("SUCCESS worker for: %v\n", p, res)
254
+ log.Debug("SUCCESS worker for: %v", p, res)
255
r.Lock()
256
r.result = res
257
r.Unlock()
@@ -235,9 +263,4 @@ func (r *dhtQueryRunner) queryPeer(p peer.Peer) {
263
r.addPeerToQuery(next, p)
264
}
265
}
238
-
239
- // signal we're done proccessing peer p
240
- log.Debug("completing worker for: %v\n", p)
241
- r.peersRemaining.Decrement(1)
242
- r.rateLimit <- struct{}{}
266
}
routing/dht/routing.go
+3
-3
@@ -29,7 +29,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
29
peers = append(peers, npeers...)
30
}
31
32
- query := newQuery(key, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
32
+ query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
33
log.Debug("%s PutValue qry part %v", dht.self, p)
34
err := dht.putValueToNetwork(ctx, p, string(key), value)
35
if err != nil {
@@ -65,7 +65,7 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
65
}
66
67
// setup the Query
68
- query := newQuery(key, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
68
+ query := newQuery(key, dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
69
70
val, peers, err := dht.getValueOrPeers(ctx, p, key, routeLevel)
71
if err != nil {
@@ -230,7 +230,7 @@ func (dht *IpfsDHT) findPeerMultiple(ctx context.Context, id peer.ID) (peer.Peer
230
}
231
232
// setup query function
233
- query := newQuery(u.Key(id), func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
233
+ query := newQuery(u.Key(id), dht.network, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
234
pmes, err := dht.findPeerSingle(ctx, p, id, routeLevel)
235
if err != nil {
236
log.Error("%s getPeer error: %v", dht.self, err)