@cryptotaxi247 / kubo / commits / a529378ce

fix(bitswap) don't 'go' local function calls

Brian Tiger Chow committed Nov 5, 2014 at 07:05 UTC a529378ce075c262a7daca6f56ab7f6f52aa6e6a
1 file changed +26 -15
exchange/bitswap/bitswap.go
+26 -15
@@ -21,15 +21,28 @@ import (
21 var log = u.Logger("bitswap")
22
23 // NetMessageSession initializes a BitSwap session that communicates over the
24 -// provided NetMessage service
25 -func NetMessageSession(parent context.Context, p peer.Peer,
24 +// provided NetMessage service.
25 +// Runs until context is cancelled
26 +func NetMessageSession(ctx context.Context, p peer.Peer,
27 net inet.Network, srv inet.Service, directory bsnet.Routing,
28 d ds.ThreadSafeDatastore, nice bool) exchange.Interface {
29
30 networkAdapter := bsnet.NetMessageAdapter(srv, net, nil)
31 +
32 + notif := notifications.New()
33 +
34 + go func() {
35 + for {
36 + select {
37 + case <-ctx.Done():
38 + notif.Shutdown()
39 + }
40 + }
41 + }()
42 +
43 bs := &bitswap{
44 blockstore: blockstore.NewBlockstore(d),
32 - notifications: notifications.New(), // TODO Shutdown()
45 + notifications: notif,
46 strategy: strategy.New(nice),
47 routing: directory,
48 sender: networkAdapter,
@@ -119,15 +132,14 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
132 case block := <-promise:
133 cancelFunc()
134 bs.wantlist.Remove(k)
122 - // TODO remove from wantlist
135 return &block, nil
136 case <-parent.Done():
137 return nil, parent.Err()
138 }
139 }
140
129 -// HasBlock announces the existance of a block to bitswap, potentially sending
130 -// it to peers (Partners) whose WantLists include it.
141 +// HasBlock announces the existance of a block to this bitswap service. The
142 +// service will potentially notify its peers.
143 func (bs *bitswap) HasBlock(ctx context.Context, blk blocks.Block) error {
144 log.Debugf("Has Block %v", blk.Key())
145 bs.wantlist.Remove(blk.Key())
@@ -162,13 +174,11 @@ func (bs *bitswap) ReceiveMessage(ctx context.Context, p peer.Peer, incoming bsm
174 if err := bs.blockstore.Put(&block); err != nil {
175 continue // FIXME(brian): err ignored
176 }
165 - go bs.notifications.Publish(block)
166 - go func(block blocks.Block) {
167 - err := bs.HasBlock(ctx, block) // FIXME err ignored
168 - if err != nil {
169 - log.Warningf("HasBlock errored: %s", err)
170 - }
171 - }(block)
177 + bs.notifications.Publish(block)
178 + err := bs.HasBlock(ctx, block)
179 + if err != nil {
180 + log.Warningf("HasBlock errored: %s", err)
181 + }
182 }
183
184 message := bsmsg.New()
@@ -202,11 +212,12 @@ func (bs *bitswap) ReceiveError(err error) {
212 // sent
213 func (bs *bitswap) send(ctx context.Context, p peer.Peer, m bsmsg.BitSwapMessage) {
214 bs.sender.SendMessage(ctx, p, m)
205 - go bs.strategy.MessageSent(p, m)
215 + bs.strategy.MessageSent(p, m)
216 }
217
218 func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block) {
219 log.Debugf("Sending %v to peers that want it", block.Key())
220 +
221 for _, p := range bs.strategy.Peers() {
222 if bs.strategy.BlockIsWantedByPeer(block.Key(), p) {
223 log.Debugf("%v wants %v", p, block.Key())
@@ -216,7 +227,7 @@ func (bs *bitswap) sendToPeersThatWant(ctx context.Context, block blocks.Block)
227 for _, wanted := range bs.wantlist.Keys() {
228 message.AddWanted(wanted)
229 }
219 - go bs.send(ctx, p, message)
230 + bs.send(ctx, p, message)
231 }
232 }
233 }