@cryptotaxi247 / kubo / commits / c7c085970

fix(exchange) allow exchange to be closed

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>

Brian Tiger Chow committed Nov 20, 2014 at 18:34 UTC c7c085970ef325099fd7398a18f78b837c728f39
3 files changed +22 -4
exchange/bitswap/bitswap.go
+12 -1
@@ -25,9 +25,11 @@ var log = eventlog.Logger("bitswap")
25 // provided BitSwapNetwork. This function registers the returned instance as
26 // the network delegate.
27 // Runs until context is cancelled
28 -func New(ctx context.Context, p peer.Peer, network bsnet.BitSwapNetwork, routing bsnet.Routing,
28 +func New(parent context.Context, p peer.Peer, network bsnet.BitSwapNetwork, routing bsnet.Routing,
29 bstore blockstore.Blockstore, nice bool) exchange.Interface {
30
31 + ctx, cancelFunc := context.WithCancel(parent)
32 +
33 notif := notifications.New()
34 go func() {
35 <-ctx.Done()
@@ -36,6 +38,7 @@ func New(ctx context.Context, p peer.Peer, network bsnet.BitSwapNetwork, routing
38
39 bs := &bitswap{
40 blockstore: bstore,
41 + cancelFunc: cancelFunc,
42 notifications: notif,
43 strategy: strategy.New(nice),
44 routing: routing,
@@ -75,6 +78,9 @@ type bitswap struct {
78 strategy strategy.Strategy
79
80 wantlist u.KeySet
81 +
82 + // cancelFunc signals cancellation to the bitswap event loop
83 + cancelFunc func()
84 }
85
86 // GetBlock attempts to retrieve a particular block from peers within the
@@ -295,3 +301,8 @@ func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block)
301 }
302 }
303 }
304 +
305 +func (bs *bitswap) Close() error {
306 + bs.cancelFunc()
307 + return nil // to conform to Closer interface
308 +}
exchange/interface.go
+4 -1
@@ -2,6 +2,8 @@
2 package exchange
3
4 import (
5 + "io"
6 +
7 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
9 blocks "github.com/jbenet/go-ipfs/blocks"
@@ -11,11 +13,12 @@ import (
13 // Any type that implements exchange.Interface may be used as an IPFS block
14 // exchange protocol.
15 type Interface interface {
14 -
16 // GetBlock returns the block associated with a given key.
17 GetBlock(context.Context, u.Key) (*blocks.Block, error)
18
19 // TODO Should callers be concerned with whether the block was made
20 // available on the network?
21 HasBlock(context.Context, blocks.Block) error
22 +
23 + io.Closer
24 }
exchange/offline/offline.go
+6 -2
@@ -20,8 +20,7 @@ func Exchange() exchange.Interface {
20
21 // offlineExchange implements the Exchange interface but doesn't return blocks.
22 // For use in offline mode.
23 -type offlineExchange struct {
24 -}
23 +type offlineExchange struct{}
24
25 // GetBlock returns nil to signal that a block could not be retrieved for the
26 // given key.
@@ -34,3 +33,8 @@ func (_ *offlineExchange) GetBlock(context.Context, u.Key) (*blocks.Block, error
33 func (_ *offlineExchange) HasBlock(context.Context, blocks.Block) error {
34 return nil
35 }
36 +
37 +// Close always returns nil.
38 +func (_ *offlineExchange) Close() error {
39 + return nil
40 +}