@cryptotaxi247 / kubo / commits / 800a3eb94

refactor(bitswap) move workers to bottom of file

Brian Tiger Chow committed Jan 31, 2015 at 01:41 UTC 800a3eb940f7b3e9d9264419782a9b66e6e04f0c
1 file changed +67 -67
exchange/bitswap/bitswap.go
+67 -67
@@ -262,73 +262,6 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context, entries []wantli
262 }
263 }
264
265 -func (bs *bitswap) taskWorker(ctx context.Context) {
266 - defer log.Info("bitswap task worker shutting down...")
267 - for {
268 - select {
269 - case <-ctx.Done():
270 - return
271 - case nextEnvelope := <-bs.engine.Outbox():
272 - select {
273 - case <-ctx.Done():
274 - return
275 - case envelope, ok := <-nextEnvelope:
276 - if !ok {
277 - continue
278 - }
279 - log.Event(ctx, "deliverBlocks", envelope.Message, envelope.Peer)
280 - bs.send(ctx, envelope.Peer, envelope.Message)
281 - }
282 - }
283 - }
284 -}
285 -
286 -// TODO ensure only one active request per key
287 -func (bs *bitswap) clientWorker(parent context.Context) {
288 -
289 - defer log.Info("bitswap client worker shutting down...")
290 -
291 - ctx, cancel := context.WithCancel(parent)
292 -
293 - broadcastSignal := time.After(rebroadcastDelay.Get())
294 - defer cancel()
295 -
296 - for {
297 - select {
298 - case <-time.Tick(10 * time.Second):
299 - n := bs.wantlist.Len()
300 - if n > 0 {
301 - log.Debug(n, inflect.FromNumber("keys", n), "in bitswap wantlist")
302 - }
303 - case <-broadcastSignal: // resend unfulfilled wantlist keys
304 - entries := bs.wantlist.Entries()
305 - if len(entries) > 0 {
306 - bs.sendWantlistToProviders(ctx, entries)
307 - }
308 - broadcastSignal = time.After(rebroadcastDelay.Get())
309 - case keys := <-bs.batchRequests:
310 - if len(keys) == 0 {
311 - log.Warning("Received batch request for zero blocks")
312 - continue
313 - }
314 - for i, k := range keys {
315 - bs.wantlist.Add(k, kMaxPriority-i)
316 - }
317 - // NB: Optimization. Assumes that providers of key[0] are likely to
318 - // be able to provide for all keys. This currently holds true in most
319 - // every situation. Later, this assumption may not hold as true.
320 - child, _ := context.WithTimeout(ctx, providerRequestTimeout)
321 - providers := bs.network.FindProvidersAsync(child, keys[0], maxProvidersPerRequest)
322 - err := bs.sendWantlistToPeers(ctx, providers)
323 - if err != nil {
324 - log.Errorf("error sending wantlist: %s", err)
325 - }
326 - case <-parent.Done():
327 - return
328 - }
329 - }
330 -}
331 -
265 // TODO(brian): handle errors
266 func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.ID, incoming bsmsg.BitSwapMessage) (
267 peer.ID, bsmsg.BitSwapMessage) {
@@ -419,3 +352,70 @@ func (bs *bitswap) send(ctx context.Context, p peer.ID, m bsmsg.BitSwapMessage)
352 func (bs *bitswap) Close() error {
353 return bs.process.Close()
354 }
355 +
356 +func (bs *bitswap) taskWorker(ctx context.Context) {
357 + defer log.Info("bitswap task worker shutting down...")
358 + for {
359 + select {
360 + case <-ctx.Done():
361 + return
362 + case nextEnvelope := <-bs.engine.Outbox():
363 + select {
364 + case <-ctx.Done():
365 + return
366 + case envelope, ok := <-nextEnvelope:
367 + if !ok {
368 + continue
369 + }
370 + log.Event(ctx, "deliverBlocks", envelope.Message, envelope.Peer)
371 + bs.send(ctx, envelope.Peer, envelope.Message)
372 + }
373 + }
374 + }
375 +}
376 +
377 +// TODO ensure only one active request per key
378 +func (bs *bitswap) clientWorker(parent context.Context) {
379 +
380 + defer log.Info("bitswap client worker shutting down...")
381 +
382 + ctx, cancel := context.WithCancel(parent)
383 +
384 + broadcastSignal := time.After(rebroadcastDelay.Get())
385 + defer cancel()
386 +
387 + for {
388 + select {
389 + case <-time.Tick(10 * time.Second):
390 + n := bs.wantlist.Len()
391 + if n > 0 {
392 + log.Debug(n, inflect.FromNumber("keys", n), "in bitswap wantlist")
393 + }
394 + case <-broadcastSignal: // resend unfulfilled wantlist keys
395 + entries := bs.wantlist.Entries()
396 + if len(entries) > 0 {
397 + bs.sendWantlistToProviders(ctx, entries)
398 + }
399 + broadcastSignal = time.After(rebroadcastDelay.Get())
400 + case keys := <-bs.batchRequests:
401 + if len(keys) == 0 {
402 + log.Warning("Received batch request for zero blocks")
403 + continue
404 + }
405 + for i, k := range keys {
406 + bs.wantlist.Add(k, kMaxPriority-i)
407 + }
408 + // NB: Optimization. Assumes that providers of key[0] are likely to
409 + // be able to provide for all keys. This currently holds true in most
410 + // every situation. Later, this assumption may not hold as true.
411 + child, _ := context.WithTimeout(ctx, providerRequestTimeout)
412 + providers := bs.network.FindProvidersAsync(child, keys[0], maxProvidersPerRequest)
413 + err := bs.sendWantlistToPeers(ctx, providers)
414 + if err != nil {
415 + log.Errorf("error sending wantlist: %s", err)
416 + }
417 + case <-parent.Done():
418 + return
419 + }
420 + }
421 +}