@cryptotaxi247 / kubo / commits / 6ab4bfea9

turn tests down a bit and better context passing

Jeromy committed May 16, 2015 at 17:46 UTC 6ab4bfea95464a3e995425a716213ed342dbeac5
3 files changed +20 -14
exchange/bitswap/bitswap.go
+2 -2
@@ -86,9 +86,9 @@ func New(parent context.Context, p peer.ID, network bsnet.BitSwapNetwork,
86 process: px,
87 newBlocks: make(chan *blocks.Block, HasBlockBufferSize),
88 provideKeys: make(chan u.Key),
89 - wm: NewWantManager(network),
89 + wm: NewWantManager(ctx, network),
90 }
91 - go bs.wm.Run(ctx)
91 + go bs.wm.Run()
92 network.SetDelegate(bs)
93
94 // Start up bitswaps async worker routines
exchange/bitswap/bitswap_test.go
+2 -2
@@ -92,7 +92,7 @@ func TestLargeSwarm(t *testing.T) {
92 if testing.Short() {
93 t.SkipNow()
94 }
95 - numInstances := 500
95 + numInstances := 100
96 numBlocks := 2
97 if detectrace.WithRace() {
98 // when running with the race detector, 500 instances launches
@@ -124,7 +124,6 @@ func TestLargeFileTwoPeers(t *testing.T) {
124 if testing.Short() {
125 t.SkipNow()
126 }
127 - t.Parallel()
127 numInstances := 2
128 numBlocks := 100
129 PerformDistributionTest(t, numInstances, numBlocks)
@@ -164,6 +163,7 @@ func PerformDistributionTest(t *testing.T, numInstances, numBlocks int) {
163 }
164 for _ = range outch {
165 }
166 + log.Error("DONE")
167 }(inst)
168 }
169 wg.Wait()
exchange/bitswap/wantmanager.go renamed
+16 -10
@@ -28,9 +28,11 @@ type WantManager struct {
28 wl *wantlist.Wantlist
29
30 network bsnet.BitSwapNetwork
31 +
32 + ctx context.Context
33 }
34
33 -func NewWantManager(network bsnet.BitSwapNetwork) *WantManager {
35 +func NewWantManager(ctx context.Context, network bsnet.BitSwapNetwork) *WantManager {
36 return &WantManager{
37 incoming: make(chan []*bsmsg.Entry, 10),
38 connect: make(chan peer.ID, 10),
@@ -38,6 +40,7 @@ func NewWantManager(network bsnet.BitSwapNetwork) *WantManager {
40 peers: make(map[peer.ID]*msgQueue),
41 wl: wantlist.New(),
42 network: network,
43 + ctx: ctx,
44 }
45 }
46
@@ -80,7 +83,10 @@ func (pm *WantManager) addEntries(ks []u.Key, cancel bool) {
83 },
84 })
85 }
83 - pm.incoming <- entries
86 + select {
87 + case pm.incoming <- entries:
88 + case <-pm.ctx.Done():
89 + }
90 }
91
92 func (pm *WantManager) SendBlock(ctx context.Context, env *engine.Envelope) {
@@ -97,7 +103,7 @@ func (pm *WantManager) SendBlock(ctx context.Context, env *engine.Envelope) {
103 }
104 }
105
100 -func (pm *WantManager) startPeerHandler(ctx context.Context, p peer.ID) *msgQueue {
106 +func (pm *WantManager) startPeerHandler(p peer.ID) *msgQueue {
107 _, ok := pm.peers[p]
108 if ok {
109 // TODO: log an error?
@@ -116,7 +122,7 @@ func (pm *WantManager) startPeerHandler(ctx context.Context, p peer.ID) *msgQueu
122 mq.work <- struct{}{}
123
124 pm.peers[p] = mq
119 - go pm.runQueue(ctx, mq)
125 + go pm.runQueue(mq)
126 return mq
127 }
128
@@ -131,12 +137,12 @@ func (pm *WantManager) stopPeerHandler(p peer.ID) {
137 delete(pm.peers, p)
138 }
139
134 -func (pm *WantManager) runQueue(ctx context.Context, mq *msgQueue) {
140 +func (pm *WantManager) runQueue(mq *msgQueue) {
141 for {
142 select {
143 case <-mq.work: // there is work to be done
144
139 - err := pm.network.ConnectTo(ctx, mq.p)
145 + err := pm.network.ConnectTo(pm.ctx, mq.p)
146 if err != nil {
147 log.Error(err)
148 // TODO: cant connect, what now?
@@ -153,7 +159,7 @@ func (pm *WantManager) runQueue(ctx context.Context, mq *msgQueue) {
159 mq.outlk.Unlock()
160
161 // send wantlist updates
156 - err = pm.network.SendMessage(ctx, mq.p, wlm)
162 + err = pm.network.SendMessage(pm.ctx, mq.p, wlm)
163 if err != nil {
164 log.Error("bitswap send error: ", err)
165 // TODO: what do we do if this fails?
@@ -173,7 +179,7 @@ func (pm *WantManager) Disconnected(p peer.ID) {
179 }
180
181 // TODO: use goprocess here once i trust it
176 -func (pm *WantManager) Run(ctx context.Context) {
182 +func (pm *WantManager) Run() {
183 for {
184 select {
185 case entries := <-pm.incoming:
@@ -193,10 +199,10 @@ func (pm *WantManager) Run(ctx context.Context) {
199 }
200
201 case p := <-pm.connect:
196 - pm.startPeerHandler(ctx, p)
202 + pm.startPeerHandler(p)
203 case p := <-pm.disconnect:
204 pm.stopPeerHandler(p)
199 - case <-ctx.Done():
205 + case <-pm.ctx.Done():
206 return
207 }
208 }