respect contexts in a more timely manner
Jeromy committed
Mar 8, 2015 at 14:10 UTC
8ed0f4b854272b5369175ae2e3d9979cf11ddf54
1 file changed
+39
-12
exchange/bitswap/bitswap.go
+39
-12
@@ -227,21 +227,40 @@ func (bs *Bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
227
func (bs *Bitswap) sendWantlistMsgToPeers(ctx context.Context, m bsmsg.BitSwapMessage, peers <-chan peer.ID) error {
228
set := pset.New()
229
wg := sync.WaitGroup{}
230
- for peerToQuery := range peers {
230
232
- if !set.TryAdd(peerToQuery) { //Do once per peer
233
- continue
234
- }
231
+loop:
232
+ for {
233
+ select {
234
+ case peerToQuery, ok := <-peers:
235
+ if !ok {
236
+ break loop
237
+ }
238
236
- wg.Add(1)
237
- go func(p peer.ID) {
238
- defer wg.Done()
239
- if err := bs.send(ctx, p, m); err != nil {
240
- log.Debug(err) // TODO remove if too verbose
239
+ if !set.TryAdd(peerToQuery) { //Do once per peer
240
+ continue
241
}
242
- }(peerToQuery)
242
+
243
+ wg.Add(1)
244
+ go func(p peer.ID) {
245
+ defer wg.Done()
246
+ if err := bs.send(ctx, p, m); err != nil {
247
+ log.Debug(err) // TODO remove if too verbose
248
+ }
249
+ }(peerToQuery)
250
+ case <-ctx.Done():
251
+ return nil
252
+ }
253
+ }
254
+ done := make(chan struct{})
255
+ go func() {
256
+ wg.Wait()
257
+ close(done)
258
+ }()
259
+
260
+ select {
261
+ case <-done:
262
+ case <-ctx.Done():
263
}
244
- wg.Wait()
264
return nil
265
}
266
@@ -385,7 +404,15 @@ func (bs *Bitswap) wantNewBlocks(ctx context.Context, bkeys []u.Key) {
404
}
405
}(p)
406
}
388
- wg.Wait()
407
+ done := make(chan struct{})
408
+ go func() {
409
+ wg.Wait()
410
+ close(done)
411
+ }()
412
+ select {
413
+ case <-done:
414
+ case <-ctx.Done():
415
+ }
416
}
417
418
func (bs *Bitswap) ReceiveError(err error) {