fix closing and removal of sessions
License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>
Jeromy committed
Jul 10, 2017 at 23:05 UTC
8be07cabd01e55b2881d0f7ae43e22e5f1e12f0e
2 files changed
+67
-2
exchange/bitswap/session.go
+24
-2
@@ -78,13 +78,28 @@ func (bs *Bitswap) NewSession(ctx context.Context) *Session {
78
return s
79
}
80
81
+func (bs *Bitswap) removeSession(s *Session) {
82
+ bs.sessLk.Lock()
83
+ defer bs.sessLk.Unlock()
84
+ for i := 0; i < len(bs.sessions); i++ {
85
+ if bs.sessions[i] == s {
86
+ bs.sessions[i] = bs.sessions[len(bs.sessions)-1]
87
+ bs.sessions = bs.sessions[:len(bs.sessions)-1]
88
+ return
89
+ }
90
+ }
91
+}
92
+
93
type blkRecv struct {
94
from peer.ID
95
blk blocks.Block
96
}
97
98
func (s *Session) receiveBlockFrom(from peer.ID, blk blocks.Block) {
87
- s.incoming <- blkRecv{from: from, blk: blk}
99
+ select {
100
+ case s.incoming <- blkRecv{from: from, blk: blk}:
101
+ case <-s.ctx.Done():
102
+ }
103
}
104
105
type interestReq struct {
@@ -105,7 +120,13 @@ func (s *Session) isLiveWant(c *cid.Cid) bool {
120
c: c,
121
resp: resp,
122
}
108
- return <-resp
123
+
124
+ select {
125
+ case want := <-resp:
126
+ return want
127
+ case <-s.ctx.Done():
128
+ return false
129
+ }
130
}
131
132
func (s *Session) interestedIn(c *cid.Cid) bool {
@@ -194,6 +215,7 @@ func (s *Session) run(ctx context.Context) {
215
lwchk.resp <- s.cidIsWanted(lwchk.c)
216
case <-ctx.Done():
217
s.tick.Stop()
218
+ s.bs.removeSession(s)
219
return
220
}
221
}
exchange/bitswap/session_test.go
+43
@@ -242,3 +242,46 @@ func TestPutAfterSessionCacheEvict(t *testing.T) {
242
t.Fatal("timed out waiting for block")
243
}
244
}
245
+
246
+func TestMultipleSessions(t *testing.T) {
247
+ ctx, cancel := context.WithCancel(context.Background())
248
+ defer cancel()
249
+
250
+ vnet := getVirtualNetwork()
251
+ sesgen := NewTestSessionGenerator(vnet)
252
+ defer sesgen.Close()
253
+ bgen := blocksutil.NewBlockGenerator()
254
+
255
+ blk := bgen.Blocks(1)[0]
256
+ inst := sesgen.Instances(2)
257
+
258
+ a := inst[0]
259
+ b := inst[1]
260
+
261
+ ctx1, cancel1 := context.WithCancel(ctx)
262
+ ses := a.Exchange.NewSession(ctx1)
263
+
264
+ blkch, err := ses.GetBlocks(ctx, []*cid.Cid{blk.Cid()})
265
+ if err != nil {
266
+ t.Fatal(err)
267
+ }
268
+ cancel1()
269
+
270
+ ses2 := a.Exchange.NewSession(ctx)
271
+ blkch2, err := ses2.GetBlocks(ctx, []*cid.Cid{blk.Cid()})
272
+ if err != nil {
273
+ t.Fatal(err)
274
+ }
275
+
276
+ time.Sleep(time.Millisecond * 10)
277
+ if err := b.Exchange.HasBlock(blk); err != nil {
278
+ t.Fatal(err)
279
+ }
280
+
281
+ select {
282
+ case <-blkch2:
283
+ case <-time.After(time.Second * 20):
284
+ t.Fatal("bad juju")
285
+ }
286
+ _ = blkch
287
+}