@cryptotaxi247 / kubo / commits / ffba03146

test closing/cancellation

- does end properly - no goroutines leaked!

Juan Batiz-Benet committed Oct 18, 2014 at 03:27 UTC ffba031469f5982b33a0a740db1484dc9b04e024
3 files changed +184 -4
crypto/spipe/handshake.go
+22 -3
@@ -18,6 +18,7 @@ import (
18 "hash"
19
20 proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
21 +
22 ci "github.com/jbenet/go-ipfs/crypto"
23 peer "github.com/jbenet/go-ipfs/peer"
24 u "github.com/jbenet/go-ipfs/util"
@@ -229,7 +230,15 @@ func (s *SecurePipe) handleSecureIn(hashType string, tIV, tCKey, tMKey []byte) {
230 theirMac, macSize := makeMac(hashType, tMKey)
231
232 for {
232 - data, ok := <-s.insecure.In
233 + var data []byte
234 + ok := true
235 +
236 + select {
237 + case <-s.ctx.Done():
238 + ok = false // return out
239 + case data, ok = <-s.insecure.In:
240 + }
241 +
242 if !ok {
243 close(s.Duplex.In)
244 return
@@ -266,8 +275,17 @@ func (s *SecurePipe) handleSecureOut(hashType string, mIV, mCKey, mMKey []byte)
275 myMac, macSize := makeMac(hashType, mMKey)
276
277 for {
269 - data, ok := <-s.Out
278 + var data []byte
279 + ok := true
280 +
281 + select {
282 + case <-s.ctx.Done():
283 + ok = false // return out
284 + case data, ok = <-s.Out:
285 + }
286 +
287 if !ok {
288 + close(s.insecure.Out)
289 return
290 }
291
@@ -363,7 +381,8 @@ func getOrConstructPeer(peers peer.Peerstore, rpk ci.PubKey) (*peer.Peer, error)
381 // this shouldn't ever happen, given we hashed, etc, but it could mean
382 // expected code (or protocol) invariants violated.
383 if !npeer.PubKey.Equals(rpk) {
366 - return nil, fmt.Errorf("WARNING: PubKey mismatch: %v", npeer)
384 + log.Error("WARNING: PubKey mismatch: %v", npeer)
385 + panic("secure channel pubkey mismatch")
386 }
387 return npeer, nil
388 }
net/conn/conn_test.go
+159
@@ -1,7 +1,13 @@
1 package conn
2
3 import (
4 + "bytes"
5 + "fmt"
6 + "runtime"
7 + "strconv"
8 + "sync"
9 "testing"
10 + "time"
11
12 ci "github.com/jbenet/go-ipfs/crypto"
13 peer "github.com/jbenet/go-ipfs/peer"
@@ -55,6 +61,43 @@ func echo(ctx context.Context, c Conn) {
61 }
62 }
63
64 +func setupConn(t *testing.T, ctx context.Context, a1, a2 string) (a, b Conn) {
65 +
66 + p1, err := setupPeer(a1)
67 + if err != nil {
68 + t.Fatal("error setting up peer", err)
69 + }
70 +
71 + p2, err := setupPeer(a2)
72 + if err != nil {
73 + t.Fatal("error setting up peer", err)
74 + }
75 +
76 + laddr := p1.NetAddress("tcp")
77 + if laddr == nil {
78 + t.Fatal("Listen address is nil.")
79 + }
80 +
81 + l1, err := Listen(ctx, laddr, p1, peer.NewPeerstore())
82 + if err != nil {
83 + t.Fatal(err)
84 + }
85 +
86 + d2 := &Dialer{
87 + Peerstore: peer.NewPeerstore(),
88 + LocalPeer: p2,
89 + }
90 +
91 + c2, err := d2.Dial(ctx, "tcp", p1)
92 + if err != nil {
93 + t.Fatal("error dialing peer", err)
94 + }
95 +
96 + c1 := <-l1.Accept()
97 +
98 + return c1, c2
99 +}
100 +
101 func TestDialer(t *testing.T) {
102
103 p1, err := setupPeer("/ip4/127.0.0.1/tcp/1234")
@@ -113,3 +156,119 @@ func TestDialer(t *testing.T) {
156 l.Close()
157 cancel()
158 }
159 +
160 +func TestClose(t *testing.T) {
161 +
162 + ctx, cancel := context.WithCancel(context.Background())
163 + c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
164 +
165 + select {
166 + case <-c1.Done():
167 + t.Fatal("done before close")
168 + case <-c2.Done():
169 + t.Fatal("done before close")
170 + default:
171 + }
172 +
173 + c1.Close()
174 +
175 + select {
176 + case <-c1.Done():
177 + default:
178 + t.Fatal("not done after cancel")
179 + }
180 +
181 + c2.Close()
182 +
183 + select {
184 + case <-c2.Done():
185 + default:
186 + t.Fatal("not done after cancel")
187 + }
188 +
189 + cancel() // close the listener :P
190 +}
191 +
192 +func TestCancel(t *testing.T) {
193 +
194 + ctx, cancel := context.WithCancel(context.Background())
195 + c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
196 +
197 + select {
198 + case <-c1.Done():
199 + t.Fatal("done before close")
200 + case <-c2.Done():
201 + t.Fatal("done before close")
202 + default:
203 + }
204 +
205 + cancel()
206 +
207 + // wait to ensure other goroutines run and close things.
208 + <-time.After(time.Microsecond * 10)
209 + // test that cancel called Close.
210 +
211 + select {
212 + case <-c1.Done():
213 + default:
214 + t.Fatal("not done after cancel")
215 + }
216 +
217 + select {
218 + case <-c2.Done():
219 + default:
220 + t.Fatal("not done after cancel")
221 + }
222 +
223 +}
224 +
225 +func TestCloseLeak(t *testing.T) {
226 +
227 + var wg sync.WaitGroup
228 +
229 + runPair := func(p1, p2, num int) {
230 + a1 := strconv.Itoa(p1)
231 + a2 := strconv.Itoa(p2)
232 + ctx, cancel := context.WithCancel(context.Background())
233 + c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/"+a1, "/ip4/127.0.0.1/tcp/"+a2)
234 +
235 + for i := 0; i < num; i++ {
236 + b1 := []byte("beep")
237 + c1.Out() <- b1
238 + b2 := <-c2.In()
239 + if !bytes.Equal(b1, b2) {
240 + panic("bytes not equal")
241 + }
242 +
243 + b2 = []byte("boop")
244 + c2.Out() <- b2
245 + b1 = <-c1.In()
246 + if !bytes.Equal(b1, b2) {
247 + panic("bytes not equal")
248 + }
249 +
250 + <-time.After(time.Microsecond * 5)
251 + }
252 +
253 + cancel() // close the listener
254 + wg.Done()
255 + }
256 +
257 + var cons = 20
258 + var msgs = 100
259 + fmt.Printf("Running %d connections * %d msgs.\n", cons, msgs)
260 + for i := 0; i < cons; i++ {
261 + wg.Add(1)
262 + go runPair(2000+i, 2001+i, msgs)
263 + }
264 +
265 + fmt.Printf("Waiting...\n")
266 + wg.Wait()
267 + // done!
268 +
269 + <-time.After(time.Microsecond * 100)
270 + if runtime.NumGoroutine() > 10 {
271 + // panic("uncomment me to debug")
272 + t.Fatal("leaking goroutines:", runtime.NumGoroutine())
273 + }
274 +}
net/conn/interface.go
+3 -1
@@ -12,6 +12,8 @@ type Map map[u.Key]Conn
12
13 // Conn is a generic message-based Peer-to-Peer connection.
14 type Conn interface {
15 + // implement ContextCloser too!
16 + ContextCloser
17
18 // LocalPeer is the Peer on this side
19 LocalPeer() *peer.Peer
@@ -26,7 +28,7 @@ type Conn interface {
28 Out() chan<- []byte
29
30 // Close ends the connection
29 - Close() error
31 + // Close() error -- already in ContextCloser
32 }
33
34 // Listener is an object that can accept connections. It matches net.Listener