added some logging
Juan Batiz-Benet committed
Sep 18, 2014 at 19:41 UTC
dc0fbfd3d37af81d804f36103caa0c634fcf37c2
1 file changed
+23
-3
routing/dht/query.go
+23
-3
@@ -111,7 +111,12 @@ func (r *dhtQueryRunner) Run(peers []*peer.Peer) (*dhtQueryResult, error) {
111
r.addPeerToQuery(p, nil) // don't have access to self here...
112
}
113
114
- // wait until we're done. yep.
114
+ // go do this thing.
115
+ go r.spawnWorkers()
116
+
117
+ // so workers are working.
118
+
119
+ // wait until they're done.
120
select {
121
case <-r.peersRemaining.Done():
122
r.cancel() // ran all and nothing. cancel all outstanding workers.
@@ -158,6 +163,8 @@ func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
163
r.peersSeen[next.Key()] = next
164
r.Unlock()
165
166
+ u.POut("adding peer to query: %v\n", next.ID.Pretty())
167
+
168
// do this after unlocking to prevent possible deadlocks.
169
r.peersRemaining.Increment(1)
170
select {
@@ -166,8 +173,9 @@ func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
173
}
174
}
175
169
-func (r *dhtQueryRunner) spawnWorkers(p *peer.Peer) {
176
+func (r *dhtQueryRunner) spawnWorkers() {
177
for {
178
+
179
select {
180
case <-r.peersRemaining.Done():
181
return
@@ -175,13 +183,19 @@ func (r *dhtQueryRunner) spawnWorkers(p *peer.Peer) {
183
case <-r.ctx.Done():
184
return
185
178
- case p := <-r.peersToQuery.DeqChan:
186
+ case p, more := <-r.peersToQuery.DeqChan:
187
+ if !more {
188
+ return // channel closed.
189
+ }
190
+ u.POut("spawning worker for: %v\n", p.ID.Pretty())
191
go r.queryPeer(p)
192
}
193
}
194
}
195
196
func (r *dhtQueryRunner) queryPeer(p *peer.Peer) {
197
+ u.POut("spawned worker for: %v\n", p.ID.Pretty())
198
+
199
// make sure we rate limit concurrency.
200
select {
201
case <-r.rateLimit:
@@ -190,27 +204,33 @@ func (r *dhtQueryRunner) queryPeer(p *peer.Peer) {
204
return
205
}
206
207
+ u.POut("running worker for: %v\n", p.ID.Pretty())
208
+
209
// finally, run the query against this peer
210
res, err := r.query.qfunc(r.ctx, p)
211
212
if err != nil {
213
+ u.POut("ERROR worker for: %v %v\n", p.ID.Pretty(), err)
214
r.Lock()
215
r.errs = append(r.errs, err)
216
r.Unlock()
217
218
} else if res.success {
219
+ u.POut("SUCCESS worker for: %v\n", p.ID.Pretty(), res)
220
r.Lock()
221
r.result = res
222
r.Unlock()
223
r.cancel() // signal to everyone that we're done.
224
225
} else if res.closerPeers != nil {
226
+ u.POut("PEERS CLOSER -- worker for: %v\n", p.ID.Pretty())
227
for _, next := range res.closerPeers {
228
r.addPeerToQuery(next, p)
229
}
230
}
231
232
// signal we're done proccessing peer p
233
+ u.POut("completing worker for: %v\n", p.ID.Pretty())
234
r.peersRemaining.Decrement(1)
235
r.rateLimit <- struct{}{}
236
}