@cryptotaxi247 / kubo / commits / 3778eedff

dont spawn so many goroutines when rebroadcasting wantlist

Jeromy committed Dec 10, 2014 at 18:47 UTC 3778eedff0ec11a8edc7ac4a82f87a10e83668fb
1 file changed +51 -18
exchange/bitswap/bitswap.go
+51 -18
@@ -201,21 +201,57 @@ func (bs *bitswap) sendWantListTo(ctx context.Context, peers <-chan peer.Peer) e
201 }
202
203 func (bs *bitswap) sendWantlistToProviders(ctx context.Context, wantlist *wl.Wantlist) {
204 + provset := make(map[u.Key]peer.Peer)
205 + provcollect := make(chan peer.Peer)
206 +
207 + ctx, cancel := context.WithCancel(ctx)
208 + defer cancel()
209 +
210 wg := sync.WaitGroup{}
211 + // Get providers for all entries in wantlist (could take a while)
212 for _, e := range wantlist.Entries() {
213 wg.Add(1)
214 go func(k u.Key) {
215 child, _ := context.WithTimeout(ctx, providerRequestTimeout)
216 providers := bs.routing.FindProvidersAsync(child, k, maxProvidersPerRequest)
217
211 - err := bs.sendWantListTo(ctx, providers)
212 - if err != nil {
213 - log.Errorf("error sending wantlist: %s", err)
218 + for prov := range providers {
219 + provcollect <- prov
220 }
221 wg.Done()
222 }(e.Value)
223 }
218 - wg.Wait()
224 +
225 + // When all workers finish, close the providers channel
226 + go func() {
227 + wg.Wait()
228 + close(provcollect)
229 + }()
230 +
231 + // Filter out duplicates,
232 + // no need to send our wantlists out twice in a given time period
233 + for {
234 + select {
235 + case p, ok := <-provcollect:
236 + if !ok {
237 + break
238 + }
239 + provset[p.Key()] = p
240 + case <-ctx.Done():
241 + log.Error("Context cancelled before we got all the providers!")
242 + return
243 + }
244 + }
245 +
246 + message := bsmsg.New()
247 + message.SetFull(true)
248 + for _, e := range bs.wantlist.Entries() {
249 + message.AddEntry(e.Value, e.Priority, false)
250 + }
251 +
252 + for _, prov := range provset {
253 + bs.send(ctx, prov, message)
254 + }
255 }
256
257 func (bs *bitswap) roundWorker(ctx context.Context) {
@@ -229,22 +265,25 @@ func (bs *bitswap) roundWorker(ctx context.Context) {
265 if err != nil {
266 log.Critical("%s", err)
267 }
232 - log.Error(alloc)
233 - bs.processStrategyAllocation(ctx, alloc)
268 + err = bs.processStrategyAllocation(ctx, alloc)
269 + if err != nil {
270 + log.Critical("Error processing strategy allocation: %s", err)
271 + }
272 }
273 }
274 }
275
238 -func (bs *bitswap) processStrategyAllocation(ctx context.Context, alloc []*strategy.Task) {
276 +func (bs *bitswap) processStrategyAllocation(ctx context.Context, alloc []*strategy.Task) error {
277 for _, t := range alloc {
278 for _, block := range t.Blocks {
279 message := bsmsg.New()
280 message.AddBlock(block)
281 if err := bs.send(ctx, t.Peer, message); err != nil {
244 - log.Errorf("Message Send Failed: %s", err)
282 + return err
283 }
284 }
285 }
286 + return nil
287 }
288
289 // TODO ensure only one active request per key
@@ -252,22 +291,16 @@ func (bs *bitswap) clientWorker(parent context.Context) {
291
292 ctx, cancel := context.WithCancel(parent)
293
255 - broadcastSignal := time.NewTicker(rebroadcastDelay)
256 - defer func() {
257 - cancel() // signal to derived async functions
258 - broadcastSignal.Stop()
259 - }()
294 + broadcastSignal := time.After(rebroadcastDelay)
295 + defer cancel()
296
297 for {
298 select {
263 - case <-broadcastSignal.C:
299 + case <-broadcastSignal:
300 // Resend unfulfilled wantlist keys
301 bs.sendWantlistToProviders(ctx, bs.wantlist)
302 + broadcastSignal = time.After(rebroadcastDelay)
303 case ks := <-bs.batchRequests:
267 - // TODO: implement batching on len(ks) > X for some X
268 - // i.e. if given 20 keys, fetch first five, then next
269 - // five, and so on, so we are more likely to be able to
270 - // effectively stream the data
304 if len(ks) == 0 {
305 log.Warning("Received batch request for zero blocks")
306 continue