separated out secure conn
Juan Batiz-Benet committed
Oct 18, 2014 at 04:02 UTC
afed188d096370aec3318fc5eb0dc03ee790c314
6 files changed
+453
-213
net/conn/closer.go
+5
-2
@@ -1,6 +1,8 @@
1
package conn
2
3
import (
4
+ "errors"
5
+
6
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
)
8
@@ -62,13 +64,14 @@ func (c *contextCloser) Done() Wait {
64
func (c *contextCloser) Close() error {
65
select {
66
case <-c.Done():
65
- panic("closed twice")
67
+ // panic("closed twice")
68
+ return errors.New("closed twice")
69
default:
70
}
71
69
- c.cancel() // release anyone waiting on the context
72
err := c.closeFunc() // actually run the close logic
73
close(c.closed) // relase everyone waiting on Done
74
+ c.cancel() // release anyone waiting on the context
75
return err
76
}
77
net/conn/conn.go
+29
-66
@@ -1,7 +1,6 @@
1
package conn
2
3
import (
4
- "errors"
4
"fmt"
5
6
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -9,7 +8,6 @@ import (
8
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9
manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
10
12
- spipe "github.com/jbenet/go-ipfs/crypto/spipe"
11
peer "github.com/jbenet/go-ipfs/peer"
12
u "github.com/jbenet/go-ipfs/util"
13
)
@@ -40,22 +38,20 @@ type singleConn struct {
38
local *peer.Peer
39
remote *peer.Peer
40
maconn manet.Conn
43
-
44
- secure *spipe.SecurePipe
45
- insecure *msgioPipe
41
+ msgio *msgioPipe
42
43
ContextCloser
44
}
45
46
// newConn constructs a new connection
47
func newSingleConn(ctx context.Context, local, remote *peer.Peer,
52
- peers peer.Peerstore, maconn manet.Conn) (Conn, error) {
48
+ maconn manet.Conn) (Conn, error) {
49
50
conn := &singleConn{
55
- local: local,
56
- remote: remote,
57
- maconn: maconn,
58
- insecure: newMsgioPipe(10),
51
+ local: local,
52
+ remote: remote,
53
+ maconn: maconn,
54
+ msgio: newMsgioPipe(10),
55
}
56
57
conn.ContextCloser = NewContextCloser(ctx, conn.close)
@@ -63,65 +59,19 @@ func newSingleConn(ctx context.Context, local, remote *peer.Peer,
59
log.Info("newSingleConn: %v to %v", local, remote)
60
61
// setup the various io goroutines
66
- go conn.insecure.outgoing.WriteTo(maconn)
67
- go conn.insecure.incoming.ReadFrom(maconn, MaxMessageSize)
68
-
69
- // perform secure handshake before returning this connection.
70
- if err := conn.secureHandshake(peers); err != nil {
71
- conn.Close()
72
- return nil, err
73
- }
62
+ go conn.msgio.outgoing.WriteTo(maconn)
63
+ go conn.msgio.incoming.ReadFrom(maconn, MaxMessageSize)
64
65
return conn, nil
66
}
67
78
-// secureHandshake performs the spipe secure handshake.
79
-func (c *singleConn) secureHandshake(peers peer.Peerstore) error {
80
- if c.secure != nil {
81
- return errors.New("Conn is already secured or being secured.")
82
- }
83
-
84
- // setup a Duplex pipe for spipe
85
- insecure := spipe.Duplex{
86
- In: c.insecure.incoming.MsgChan,
87
- Out: c.insecure.outgoing.MsgChan,
88
- }
89
-
90
- // spipe performs the secure handshake, which takes multiple RTT
91
- sp, err := spipe.NewSecurePipe(c.Context(), 10, c.local, peers, insecure)
92
- if err != nil {
93
- return err
94
- }
95
-
96
- // assign it into the conn object
97
- c.secure = sp
98
-
99
- if c.remote == nil {
100
- c.remote = c.secure.RemotePeer()
101
-
102
- } else if c.remote != c.secure.RemotePeer() {
103
- // this panic is here because this would be an insidious programmer error
104
- // that we need to ensure we catch.
105
- log.Error("%v != %v", c.remote, c.secure.RemotePeer())
106
- panic("peers not being constructed correctly.")
107
- }
108
-
109
- return nil
110
-}
111
-
68
// close is the internal close function, called by ContextCloser.Close
69
func (c *singleConn) close() error {
70
log.Debug("%s closing Conn with %s", c.local, c.remote)
71
72
// close underlying connection
73
err := c.maconn.Close()
118
-
119
- // closing channels
120
- c.insecure.outgoing.Close()
121
- if c.secure != nil { // may never have gotten here.
122
- c.secure.Close()
123
- }
124
-
74
+ c.msgio.outgoing.Close()
75
return err
76
}
77
@@ -137,12 +87,12 @@ func (c *singleConn) RemotePeer() *peer.Peer {
87
88
// In returns a readable message channel
89
func (c *singleConn) In() <-chan []byte {
140
- return c.secure.In
90
+ return c.msgio.incoming.MsgChan
91
}
92
93
// Out returns a writable message channel
94
func (c *singleConn) Out() chan<- []byte {
145
- return c.secure.Out
95
+ return c.msgio.outgoing.MsgChan
96
}
97
98
// Dialer is an object that can open connections. We could have a "convenience"
@@ -186,7 +136,12 @@ func (d *Dialer) Dial(ctx context.Context, network string, remote *peer.Peer) (C
136
log.Error("Error putting peer into peerstore: %s", remote)
137
}
138
189
- return newSingleConn(ctx, d.LocalPeer, remote, d.Peerstore, maconn)
139
+ c, err := newSingleConn(ctx, d.LocalPeer, remote, maconn)
140
+ if err != nil {
141
+ return nil, err
142
+ }
143
+
144
+ return newSecureConn(ctx, c, d.Peerstore)
145
}
146
147
// listener is an object that can accept connections. It implements Listener
@@ -240,13 +195,21 @@ func (l *listener) listen() {
195
// handle is a goroutine work function that handles the handshake.
196
// it's here only so that accepting new connections can happen quickly.
197
handle := func(maconn manet.Conn) {
243
- c, err := newSingleConn(l.Context(), l.local, nil, l.peers, maconn)
198
+ defer func() { <-sem }() // release
199
+
200
+ c, err := newSingleConn(l.Context(), l.local, nil, maconn)
201
if err != nil {
202
log.Error("Error accepting connection: %v", err)
246
- } else {
247
- l.conns <- c
203
+ return
204
}
249
- <-sem // release
205
+
206
+ sc, err := newSecureConn(l.Context(), c, l.peers)
207
+ if err != nil {
208
+ log.Error("Error securing connection: %v", err)
209
+ return
210
+ }
211
+
212
+ l.conns <- sc
213
}
214
215
for {
net/conn/conn_test.go
-145
@@ -9,154 +9,9 @@ import (
9
"testing"
10
"time"
11
12
- ci "github.com/jbenet/go-ipfs/crypto"
13
- peer "github.com/jbenet/go-ipfs/peer"
14
-
12
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
16
- ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
13
)
14
19
-func setupPeer(addr string) (*peer.Peer, error) {
20
- tcp, err := ma.NewMultiaddr(addr)
21
- if err != nil {
22
- return nil, err
23
- }
24
-
25
- sk, pk, err := ci.GenerateKeyPair(ci.RSA, 512)
26
- if err != nil {
27
- return nil, err
28
- }
29
-
30
- id, err := peer.IDFromPubKey(pk)
31
- if err != nil {
32
- return nil, err
33
- }
34
-
35
- p := &peer.Peer{ID: id}
36
- p.PrivKey = sk
37
- p.PubKey = pk
38
- p.AddAddress(tcp)
39
- return p, nil
40
-}
41
-
42
-func echoListen(ctx context.Context, listener Listener) {
43
- for {
44
- select {
45
- case <-ctx.Done():
46
- return
47
- case c := <-listener.Accept():
48
- go echo(ctx, c)
49
- }
50
- }
51
-}
52
-
53
-func echo(ctx context.Context, c Conn) {
54
- for {
55
- select {
56
- case <-ctx.Done():
57
- return
58
- case m := <-c.In():
59
- c.Out() <- m
60
- }
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")
104
- if err != nil {
105
- t.Fatal("error setting up peer", err)
106
- }
107
-
108
- p2, err := setupPeer("/ip4/127.0.0.1/tcp/3456")
109
- if err != nil {
110
- t.Fatal("error setting up peer", err)
111
- }
112
-
113
- ctx, cancel := context.WithCancel(context.Background())
114
-
115
- laddr := p1.NetAddress("tcp")
116
- if laddr == nil {
117
- t.Fatal("Listen address is nil.")
118
- }
119
-
120
- l, err := Listen(ctx, laddr, p1, peer.NewPeerstore())
121
- if err != nil {
122
- t.Fatal(err)
123
- }
124
-
125
- go echoListen(ctx, l)
126
-
127
- d := &Dialer{
128
- Peerstore: peer.NewPeerstore(),
129
- LocalPeer: p2,
130
- }
131
-
132
- c, err := d.Dial(ctx, "tcp", p1)
133
- if err != nil {
134
- t.Fatal("error dialing peer", err)
135
- }
136
-
137
- // fmt.Println("sending")
138
- c.Out() <- []byte("beep")
139
- c.Out() <- []byte("boop")
140
-
141
- out := <-c.In()
142
- // fmt.Println("recving", string(out))
143
- data := string(out)
144
- if data != "beep" {
145
- t.Error("unexpected conn output", data)
146
- }
147
-
148
- out = <-c.In()
149
- data = string(out)
150
- if string(out) != "boop" {
151
- t.Error("unexpected conn output", data)
152
- }
153
-
154
- // fmt.Println("closing")
155
- c.Close()
156
- l.Close()
157
- cancel()
158
-}
159
-
15
func TestClose(t *testing.T) {
16
17
ctx, cancel := context.WithCancel(context.Background())
net/conn/dial_test.go
new
+152
@@ -0,0 +1,152 @@
1
+package conn
2
+
3
+import (
4
+ "testing"
5
+
6
+ ci "github.com/jbenet/go-ipfs/crypto"
7
+ peer "github.com/jbenet/go-ipfs/peer"
8
+
9
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
10
+ ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
11
+)
12
+
13
+func setupPeer(addr string) (*peer.Peer, error) {
14
+ tcp, err := ma.NewMultiaddr(addr)
15
+ if err != nil {
16
+ return nil, err
17
+ }
18
+
19
+ sk, pk, err := ci.GenerateKeyPair(ci.RSA, 512)
20
+ if err != nil {
21
+ return nil, err
22
+ }
23
+
24
+ id, err := peer.IDFromPubKey(pk)
25
+ if err != nil {
26
+ return nil, err
27
+ }
28
+
29
+ p := &peer.Peer{ID: id}
30
+ p.PrivKey = sk
31
+ p.PubKey = pk
32
+ p.AddAddress(tcp)
33
+ return p, nil
34
+}
35
+
36
+func echoListen(ctx context.Context, listener Listener) {
37
+ for {
38
+ select {
39
+ case <-ctx.Done():
40
+ return
41
+ case c := <-listener.Accept():
42
+ go echo(ctx, c)
43
+ }
44
+ }
45
+}
46
+
47
+func echo(ctx context.Context, c Conn) {
48
+ for {
49
+ select {
50
+ case <-ctx.Done():
51
+ return
52
+ case m := <-c.In():
53
+ c.Out() <- m
54
+ }
55
+ }
56
+}
57
+
58
+func setupConn(t *testing.T, ctx context.Context, a1, a2 string) (a, b Conn) {
59
+
60
+ p1, err := setupPeer(a1)
61
+ if err != nil {
62
+ t.Fatal("error setting up peer", err)
63
+ }
64
+
65
+ p2, err := setupPeer(a2)
66
+ if err != nil {
67
+ t.Fatal("error setting up peer", err)
68
+ }
69
+
70
+ laddr := p1.NetAddress("tcp")
71
+ if laddr == nil {
72
+ t.Fatal("Listen address is nil.")
73
+ }
74
+
75
+ l1, err := Listen(ctx, laddr, p1, peer.NewPeerstore())
76
+ if err != nil {
77
+ t.Fatal(err)
78
+ }
79
+
80
+ d2 := &Dialer{
81
+ Peerstore: peer.NewPeerstore(),
82
+ LocalPeer: p2,
83
+ }
84
+
85
+ c2, err := d2.Dial(ctx, "tcp", p1)
86
+ if err != nil {
87
+ t.Fatal("error dialing peer", err)
88
+ }
89
+
90
+ c1 := <-l1.Accept()
91
+
92
+ return c1, c2
93
+}
94
+
95
+func TestDialer(t *testing.T) {
96
+
97
+ p1, err := setupPeer("/ip4/127.0.0.1/tcp/1234")
98
+ if err != nil {
99
+ t.Fatal("error setting up peer", err)
100
+ }
101
+
102
+ p2, err := setupPeer("/ip4/127.0.0.1/tcp/3456")
103
+ if err != nil {
104
+ t.Fatal("error setting up peer", err)
105
+ }
106
+
107
+ ctx, cancel := context.WithCancel(context.Background())
108
+
109
+ laddr := p1.NetAddress("tcp")
110
+ if laddr == nil {
111
+ t.Fatal("Listen address is nil.")
112
+ }
113
+
114
+ l, err := Listen(ctx, laddr, p1, peer.NewPeerstore())
115
+ if err != nil {
116
+ t.Fatal(err)
117
+ }
118
+
119
+ go echoListen(ctx, l)
120
+
121
+ d := &Dialer{
122
+ Peerstore: peer.NewPeerstore(),
123
+ LocalPeer: p2,
124
+ }
125
+
126
+ c, err := d.Dial(ctx, "tcp", p1)
127
+ if err != nil {
128
+ t.Fatal("error dialing peer", err)
129
+ }
130
+
131
+ // fmt.Println("sending")
132
+ c.Out() <- []byte("beep")
133
+ c.Out() <- []byte("boop")
134
+
135
+ out := <-c.In()
136
+ // fmt.Println("recving", string(out))
137
+ data := string(out)
138
+ if data != "beep" {
139
+ t.Error("unexpected conn output", data)
140
+ }
141
+
142
+ out = <-c.In()
143
+ data = string(out)
144
+ if string(out) != "boop" {
145
+ t.Error("unexpected conn output", data)
146
+ }
147
+
148
+ // fmt.Println("closing")
149
+ c.Close()
150
+ l.Close()
151
+ cancel()
152
+}
net/conn/secure_conn.go
new
+113
@@ -0,0 +1,113 @@
1
+package conn
2
+
3
+import (
4
+ "errors"
5
+
6
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
+
8
+ spipe "github.com/jbenet/go-ipfs/crypto/spipe"
9
+ peer "github.com/jbenet/go-ipfs/peer"
10
+)
11
+
12
+// secureConn wraps another Conn object with an encrypted channel.
13
+type secureConn struct {
14
+
15
+ // the wrapped conn
16
+ insecure Conn
17
+
18
+ // secure pipe, wrapping insecure
19
+ secure *spipe.SecurePipe
20
+
21
+ ContextCloser
22
+}
23
+
24
+// newConn constructs a new connection
25
+func newSecureConn(ctx context.Context, insecure Conn, peers peer.Peerstore) (Conn, error) {
26
+
27
+ conn := &secureConn{
28
+ insecure: insecure,
29
+ }
30
+ conn.ContextCloser = NewContextCloser(ctx, conn.close)
31
+
32
+ log.Debug("newSecureConn: %v to %v", insecure.LocalPeer(), insecure.RemotePeer())
33
+ // perform secure handshake before returning this connection.
34
+ if err := conn.secureHandshake(peers); err != nil {
35
+ conn.Close()
36
+ return nil, err
37
+ }
38
+ log.Debug("newSecureConn: %v to %v handshake success!", insecure.LocalPeer(), insecure.RemotePeer())
39
+
40
+ return conn, nil
41
+}
42
+
43
+// secureHandshake performs the spipe secure handshake.
44
+func (c *secureConn) secureHandshake(peers peer.Peerstore) error {
45
+ if c.secure != nil {
46
+ return errors.New("Conn is already secured or being secured.")
47
+ }
48
+
49
+ // ok to panic here if this type assertion fails. Interface hack.
50
+ // when we support wrapping other Conns, we'll need to change
51
+ // spipe to do something else.
52
+ insecureSC := c.insecure.(*singleConn)
53
+
54
+ // setup a Duplex pipe for spipe
55
+ insecureD := spipe.Duplex{
56
+ In: insecureSC.msgio.incoming.MsgChan,
57
+ Out: insecureSC.msgio.outgoing.MsgChan,
58
+ }
59
+
60
+ // spipe performs the secure handshake, which takes multiple RTT
61
+ sp, err := spipe.NewSecurePipe(c.Context(), 10, c.LocalPeer(), peers, insecureD)
62
+ if err != nil {
63
+ return err
64
+ }
65
+
66
+ // assign it into the conn object
67
+ c.secure = sp
68
+
69
+ // if we do not know RemotePeer, get it from secure chan (who identifies it)
70
+ if insecureSC.remote == nil {
71
+ insecureSC.remote = c.secure.RemotePeer()
72
+
73
+ } else if insecureSC.remote != c.secure.RemotePeer() {
74
+ // this panic is here because this would be an insidious programmer error
75
+ // that we need to ensure we catch.
76
+ // update: this actually might happen under normal operation-- should
77
+ // perhaps return an error. TBD.
78
+
79
+ log.Error("secureConn peer mismatch. %v != %v", insecureSC.remote, c.secure.RemotePeer())
80
+ panic("secureConn peer mismatch. consructed incorrectly?")
81
+ }
82
+
83
+ return nil
84
+}
85
+
86
+// close is called by ContextCloser
87
+func (c *secureConn) close() error {
88
+ err := c.insecure.Close()
89
+ if c.secure != nil { // may never have gotten here.
90
+ err = c.secure.Close()
91
+ }
92
+ return err
93
+}
94
+
95
+// LocalPeer is the Peer on this side
96
+func (c *secureConn) LocalPeer() *peer.Peer {
97
+ return c.insecure.LocalPeer()
98
+}
99
+
100
+// RemotePeer is the Peer on the remote side
101
+func (c *secureConn) RemotePeer() *peer.Peer {
102
+ return c.insecure.RemotePeer()
103
+}
104
+
105
+// In returns a readable message channel
106
+func (c *secureConn) In() <-chan []byte {
107
+ return c.secure.In
108
+}
109
+
110
+// Out returns a writable message channel
111
+func (c *secureConn) Out() chan<- []byte {
112
+ return c.secure.Out
113
+}
net/conn/secure_conn_test.go
new
+154
@@ -0,0 +1,154 @@
1
+package conn
2
+
3
+import (
4
+ "bytes"
5
+ "fmt"
6
+ "runtime"
7
+ "strconv"
8
+ "sync"
9
+ "testing"
10
+ "time"
11
+
12
+ peer "github.com/jbenet/go-ipfs/peer"
13
+
14
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
15
+)
16
+
17
+func setupSecureConn(t *testing.T, c Conn) Conn {
18
+ c, ok := c.(*secureConn)
19
+ if ok {
20
+ return c
21
+ }
22
+
23
+ // shouldn't happen, because dial + listen already return secure conns.
24
+ s, err := newSecureConn(c.Context(), c, peer.NewPeerstore())
25
+ if err != nil {
26
+ t.Fatal(err)
27
+ }
28
+ return s
29
+}
30
+
31
+func TestSecureClose(t *testing.T) {
32
+
33
+ ctx, cancel := context.WithCancel(context.Background())
34
+ c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
35
+
36
+ c1 = setupSecureConn(t, c1)
37
+ c2 = setupSecureConn(t, c2)
38
+
39
+ select {
40
+ case <-c1.Done():
41
+ t.Fatal("done before close")
42
+ case <-c2.Done():
43
+ t.Fatal("done before close")
44
+ default:
45
+ }
46
+
47
+ c1.Close()
48
+
49
+ select {
50
+ case <-c1.Done():
51
+ default:
52
+ t.Fatal("not done after cancel")
53
+ }
54
+
55
+ c2.Close()
56
+
57
+ select {
58
+ case <-c2.Done():
59
+ default:
60
+ t.Fatal("not done after cancel")
61
+ }
62
+
63
+ cancel() // close the listener :P
64
+}
65
+
66
+func TestSecureCancel(t *testing.T) {
67
+
68
+ ctx, cancel := context.WithCancel(context.Background())
69
+ c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/1234", "/ip4/127.0.0.1/tcp/2345")
70
+
71
+ c1 = setupSecureConn(t, c1)
72
+ c2 = setupSecureConn(t, c2)
73
+
74
+ select {
75
+ case <-c1.Done():
76
+ t.Fatal("done before close")
77
+ case <-c2.Done():
78
+ t.Fatal("done before close")
79
+ default:
80
+ }
81
+
82
+ cancel()
83
+
84
+ // wait to ensure other goroutines run and close things.
85
+ <-time.After(time.Microsecond * 10)
86
+ // test that cancel called Close.
87
+
88
+ select {
89
+ case <-c1.Done():
90
+ default:
91
+ t.Fatal("not done after cancel")
92
+ }
93
+
94
+ select {
95
+ case <-c2.Done():
96
+ default:
97
+ t.Fatal("not done after cancel")
98
+ }
99
+
100
+}
101
+
102
+func TestSecureCloseLeak(t *testing.T) {
103
+
104
+ var wg sync.WaitGroup
105
+
106
+ runPair := func(p1, p2, num int) {
107
+ a1 := strconv.Itoa(p1)
108
+ a2 := strconv.Itoa(p2)
109
+ ctx, cancel := context.WithCancel(context.Background())
110
+ c1, c2 := setupConn(t, ctx, "/ip4/127.0.0.1/tcp/"+a1, "/ip4/127.0.0.1/tcp/"+a2)
111
+
112
+ c1 = setupSecureConn(t, c1)
113
+ c2 = setupSecureConn(t, c2)
114
+
115
+ for i := 0; i < num; i++ {
116
+ b1 := []byte("beep")
117
+ c1.Out() <- b1
118
+ b2 := <-c2.In()
119
+ if !bytes.Equal(b1, b2) {
120
+ panic("bytes not equal")
121
+ }
122
+
123
+ b2 = []byte("boop")
124
+ c2.Out() <- b2
125
+ b1 = <-c1.In()
126
+ if !bytes.Equal(b1, b2) {
127
+ panic("bytes not equal")
128
+ }
129
+
130
+ <-time.After(time.Microsecond * 5)
131
+ }
132
+
133
+ cancel() // close the listener
134
+ wg.Done()
135
+ }
136
+
137
+ var cons = 20
138
+ var msgs = 100
139
+ fmt.Printf("Running %d connections * %d msgs.\n", cons, msgs)
140
+ for i := 0; i < cons; i++ {
141
+ wg.Add(1)
142
+ go runPair(2000+i, 2001+i, msgs)
143
+ }
144
+
145
+ fmt.Printf("Waiting...\n")
146
+ wg.Wait()
147
+ // done!
148
+
149
+ <-time.After(time.Microsecond * 100)
150
+ if runtime.NumGoroutine() > 10 {
151
+ // panic("uncomment me to debug")
152
+ t.Fatal("leaking goroutines:", runtime.NumGoroutine())
153
+ }
154
+}