@cryptotaxi247 / kubo / commits / 9a2f79c42

refactor(bitswap) consolidate HasBlock

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com> Conflicts: exchange/bitswap/bitswap.go

Brian Tiger Chow committed Nov 26, 2014 at 17:51 UTC 9a2f79c42fe1714d9f854dcd5905766a65082091
1 file changed +19 -26
exchange/bitswap/bitswap.go
+19 -26
@@ -248,30 +248,19 @@ func (bs *bitswap) loop(parent context.Context) {
248 // HasBlock announces the existance of a block to this bitswap service. The
249 // service will potentially notify its peers.
250 func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
251 - // TODO check all errors
252 - log.Debugf("Has Block %s", blk.Key())
251 + if err := bs.blockstore.Put(blk); err != nil {
252 + return err
253 + }
254 bs.wantlist.Remove(blk.Key())
255 bs.notifications.Publish(blk)
255 -
256 child, _ := context.WithTimeout(ctx, hasBlockTimeout)
257 - bs.sendToPeersThatWant(child, blk)
257 + if err := bs.sendToPeersThatWant(child, blk); err != nil {
258 + return err
259 + }
260 child, _ = context.WithTimeout(ctx, hasBlockTimeout)
261 return bs.routing.Provide(child, blk.Key())
262 }
263
262 -// receiveBlock handles storing the block in the blockstore and calling HasBlock
263 -func (bs *bitswap) receiveBlock(ctx context.Context, block *blocks.Block) {
264 - // TODO verify blocks?
265 - if err := bs.blockstore.Put(block); err != nil {
266 - log.Criticalf("error putting block: %s", err)
267 - return
268 - }
269 - err := bs.HasBlock(ctx, block)
270 - if err != nil {
271 - log.Warningf("HasBlock errored: %s", err)
272 - }
273 -}
274 -
264 // TODO(brian): handle errors
265 func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsmsg.BitSwapMessage) (
266 peer.Peer, bsmsg.BitSwapMessage) {
@@ -297,7 +286,9 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
286
287 go func() {
288 for _, block := range incoming.Blocks() {
300 - bs.receiveBlock(ctx, block)
289 + if err := bs.HasBlock(ctx, block); err != nil {
290 + log.Error(err)
291 + }
292 }
293 }()
294
@@ -334,27 +325,29 @@ func (bs *bitswap) ReceiveError(err error) {
325
326 // send strives to ensure that accounting is always performed when a message is
327 // sent
337 -func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage) {
338 - bs.sender.SendMessage(ctx, p, m)
339 - bs.strategy.MessageSent(p, m)
328 +func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage) error {
329 + if err := bs.sender.SendMessage(ctx, p, m); err != nil {
330 + return err
331 + }
332 + return bs.strategy.MessageSent(p, m)
333 }
334
342 -func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block *blocks.Block) {
343 - log.Debugf("Sending %s to peers that want it", block)
344 -
335 +func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block *blocks.Block) error {
336 for _, p := range bs.strategy.Peers() {
337 if bs.strategy.BlockIsWantedByPeer(block.Key(), p) {
347 - log.Debugf("%v wants %v", p, block.Key())
338 if bs.strategy.ShouldSendBlockToPeer(block.Key(), p) {
339 message := bsmsg.New()
340 message.AddBlock(block)
341 for _, wanted := range bs.wantlist.Keys() {
342 message.AddWanted(wanted)
343 }
354 - bs.send(ctx, p, message)
344 + if err := bs.send(ctx, p, message); err != nil {
345 + return err
346 + }
347 }
348 }
349 }
350 + return nil
351 }
352
353 func (bs *bitswap) Close() error {