stress test
Juan Batiz-Benet committed
Dec 16, 2014 at 12:12 UTC
7fdafaf1e52518d21ae0008affad20c59f44bcf2
2 files changed
+107
-9
net/mock/mock.go
+11
-4
@@ -2,17 +2,21 @@
2
package mocknet
3
4
import (
5
+ "fmt"
6
"io"
7
"sync"
8
9
inet "github.com/jbenet/go-ipfs/net"
10
peer "github.com/jbenet/go-ipfs/peer"
11
+ eventlog "github.com/jbenet/go-ipfs/util/eventlog"
12
13
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14
ctxgroup "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-ctxgroup"
15
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
16
)
17
18
+var log = eventlog.Logger("mocknet")
19
+
20
type Stream struct {
21
io.Reader
22
io.Writer
@@ -98,6 +102,7 @@ func (c *Conn) removeStream(s *Stream) {
102
103
func (c *Conn) NewStreamWithProtocol(pr inet.ProtocolID, p peer.Peer) (inet.Stream, error) {
104
105
+ log.Debugf("NewStreamWithProtocol: %s --> %s", c.local, p)
106
ss, _ := newStreamPair(c.local, p)
107
108
if err := inet.WriteProtocolHeader(pr, ss); err != nil {
@@ -149,13 +154,12 @@ func MakeNetworks(ctx context.Context, peers []peer.Peer) (nets []*Network, err
154
}
155
}
156
157
+ i := 0
158
for _, n1 := range nets {
159
for _, n2 := range nets {
154
- if n1 == n2 {
155
- continue
156
- }
157
-
160
n1.conns[n2.local] = &Conn{local: n1, remote: n2}
161
+ log.Debugf("%d setup %s -> %s", i, n1, n2)
162
+ i++
163
}
164
}
165
@@ -175,6 +179,9 @@ func newNetwork(ctx context.Context, local peer.Peer, peers peer.Peerstore) (*Ne
179
n.cg.SetTeardown(n.close)
180
return n, nil
181
}
182
+func (n *Network) String() string {
183
+ return fmt.Sprintf("<Network %s - %d conns>", n.local, len(n.conns))
184
+}
185
186
func (n *Network) handle(s inet.Stream) {
187
go n.mux.Handle(s)
net/mock/mock_test.go
+96
-5
@@ -3,6 +3,8 @@ package mocknet
3
import (
4
"bytes"
5
"io"
6
+ "math/rand"
7
+ "sync"
8
"testing"
9
10
inet "github.com/jbenet/go-ipfs/net"
@@ -35,15 +37,11 @@ func TestNetworkSetup(t *testing.T) {
37
t.Error("peer mismatch")
38
}
39
38
- if len(n.conns) != (len(nets) - 1) {
40
+ if len(n.conns) != len(nets) {
41
t.Error("conn mismatch")
42
}
43
44
for _, c := range n.conns {
43
- if c.remote.local == n.local {
44
- t.Error("conn to self")
45
- }
46
-
45
if c.remote.conns[n.local] == nil {
46
t.Error("conn other side fail")
47
}
@@ -101,3 +99,96 @@ func TestStreams(t *testing.T) {
99
}
100
101
}
102
+
103
+func makePinger(st string, n int) func(inet.Stream) {
104
+ return func(s inet.Stream) {
105
+ go func() {
106
+ defer s.Close()
107
+
108
+ for i := 0; i < n; i++ {
109
+ b := make([]byte, 4+len(st))
110
+ if _, err := s.Write([]byte("ping" + st)); err != nil {
111
+ panic(err)
112
+ }
113
+ if _, err := io.ReadFull(s, b); err != nil {
114
+ panic(err)
115
+ }
116
+ if !bytes.Equal(b, []byte("pong"+st)) {
117
+ panic("bytes mismatch")
118
+ }
119
+ }
120
+ }()
121
+ }
122
+}
123
+
124
+func makePonger(st string) func(inet.Stream) {
125
+ return func(s inet.Stream) {
126
+ go func() {
127
+ defer s.Close()
128
+
129
+ for {
130
+ b := make([]byte, 4+len(st))
131
+ if _, err := io.ReadFull(s, b); err != nil {
132
+ if err == io.EOF {
133
+ return
134
+ }
135
+ panic(err)
136
+ }
137
+ if !bytes.Equal(b, []byte("ping"+st)) {
138
+ panic("bytes mismatch")
139
+ }
140
+ if _, err := s.Write([]byte("pong" + st)); err != nil {
141
+ panic(err)
142
+ }
143
+ }
144
+ }()
145
+ }
146
+}
147
+
148
+func TestStreamsStress(t *testing.T) {
149
+
150
+ peers := []peer.Peer{}
151
+ for i := 0; i < 100; i++ {
152
+ peers = append(peers, testutil.RandPeer())
153
+ }
154
+
155
+ nets, err := MakeNetworks(context.Background(), peers)
156
+ if err != nil {
157
+ t.Fatal(err)
158
+ }
159
+
160
+ protos := []inet.ProtocolID{
161
+ inet.ProtocolDHT,
162
+ inet.ProtocolBitswap,
163
+ inet.ProtocolDiag,
164
+ }
165
+
166
+ for _, n := range nets {
167
+ for _, p := range protos {
168
+ n.SetHandler(p, makePonger(string(p)))
169
+ }
170
+ }
171
+
172
+ var wg sync.WaitGroup
173
+ for i := 0; i < 1000; i++ {
174
+ wg.Add(1)
175
+ go func(i int) {
176
+ defer wg.Done()
177
+ from := rand.Intn(len(peers))
178
+ to := rand.Intn(len(peers))
179
+ p := rand.Intn(3)
180
+ proto := protos[p]
181
+ log.Debug("%d (%s) %d (%s) %d (%s)", from, nets[from], to, nets[to], p, protos[p])
182
+ s, err := nets[from].NewStream(protos[p], nets[to].local)
183
+ if err != nil {
184
+ panic(err)
185
+ }
186
+
187
+ log.Infof("%d start pinging", i)
188
+ makePinger(string(proto), rand.Intn(100))(s)
189
+ log.Infof("%d done pinging", i)
190
+ }(i)
191
+ }
192
+
193
+ wg.Done()
194
+}