document bitswap more
Jeromy committed
Nov 26, 2014 at 23:48 UTC
bc02b77b47f34586026563e241fb24914b36f90f
1 file changed
+20
-10
exchange/bitswap/bitswap.go
+20
-10
@@ -22,7 +22,8 @@ import (
22
var log = eventlog.Logger("bitswap")
23
24
// Number of providers to request for sending a wantlist to
25
-const maxProvidersPerRequest = 6
25
+// TODO: if a 'non-nice' strategy is implemented, consider increasing this value
26
+const maxProvidersPerRequest = 3
27
28
// New initializes a BitSwap instance that communicates over the
29
// provided BitSwapNetwork. This function registers the returned instance as
@@ -211,6 +212,7 @@ func (bs *bitswap) loop(parent context.Context) {
212
for {
213
select {
214
case <-broadcastSignal.C:
215
+ // Resend unfulfilled wantlist keys
216
bs.sendWantlistToProviders(ctx, bs.wantlist.Keys())
217
case ks := <-bs.batchRequests:
218
// TODO: implement batching on len(ks) > X for some X
@@ -224,6 +226,13 @@ func (bs *bitswap) loop(parent context.Context) {
226
for _, k := range ks {
227
bs.wantlist.Add(k)
228
}
229
+ // NB: send want list to providers for the first peer in this list.
230
+ // the assumption is made that the providers of the first key in
231
+ // the set are likely to have others as well.
232
+ // This currently holds true in most every situation, since when
233
+ // pinning a file, you store and provide all blocks associated with
234
+ // it. Later, this assumption may not hold as true if we implement
235
+ // newer bitswap strategies.
236
providers := bs.routing.FindProvidersAsync(ctx, ks[0], maxProvidersPerRequest)
237
238
err := bs.sendWantListTo(ctx, providers)
@@ -263,7 +272,6 @@ func (bs *bitswap) receiveBlock(ctx context.Context, block *blocks.Block) {
272
func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsmsg.BitSwapMessage) (
273
peer.Peer, bsmsg.BitSwapMessage) {
274
log.Debugf("ReceiveMessage from %s", p)
266
- log.Debugf("Message wantlist: %v", incoming.Wantlist())
275
276
if p == nil {
277
log.Error("Received message from nil peer!")
@@ -279,15 +287,17 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
287
// Record message bytes in ledger
288
// TODO: this is bad, and could be easily abused.
289
// Should only track *useful* messages in ledger
282
- bs.strategy.MessageReceived(p, incoming) // FIRST
290
+ // This call records changes to wantlists, blocks received,
291
+ // and number of bytes transfered.
292
+ bs.strategy.MessageReceived(p, incoming)
293
284
- for _, block := range incoming.Blocks() {
285
- go bs.receiveBlock(ctx, block)
286
- }
294
+ go func() {
295
+ for _, block := range incoming.Blocks() {
296
+ bs.receiveBlock(ctx, block)
297
+ }
298
+ }()
299
300
for _, key := range incoming.Wantlist() {
289
- // TODO: might be better to check if we have the block before checking
290
- // if we should send it to someone
301
if bs.strategy.ShouldSendBlockToPeer(key, p) {
302
if block, errBlockNotFound := bs.blockstore.Get(key); errBlockNotFound != nil {
303
continue
@@ -303,12 +313,12 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
313
}
314
315
blkmsg.AddBlock(block)
306
- bs.strategy.MessageSent(p, blkmsg)
316
bs.send(ctx, p, blkmsg)
317
}
318
}
319
}
320
321
+ // TODO: consider changing this function to not return anything
322
return nil, nil
323
}
324
@@ -326,7 +336,7 @@ func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage
336
}
337
338
func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block *blocks.Block) {
329
- log.Debugf("Sending %v to peers that want it", block.Key())
339
+ log.Debugf("Sending %s to peers that want it", block)
340
341
for _, p := range bs.strategy.Peers() {
342
if bs.strategy.BlockIsWantedByPeer(block.Key(), p) {