@cryptotaxi247 / kubo / commits / 58510640f

rm old mock

Juan Batiz-Benet committed Dec 17, 2014 at 07:53 UTC 58510640fa18ad793780c60f88d4ae8340155f26
2 files changed -504
net/mock/mock.go deleted
-310
@@ -1,310 +0,0 @@
1 -// Package mocknet provides a mock net.Network to test with.
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
23 - conn *Conn
24 -}
25 -
26 -func (s *Stream) Close() error {
27 - s.conn.removeStream(s)
28 - if r, ok := (s.Reader).(io.Closer); ok {
29 - r.Close()
30 - }
31 - if w, ok := (s.Writer).(io.Closer); ok {
32 - return w.Close()
33 - }
34 - return nil
35 -}
36 -
37 -func (s *Stream) Conn() inet.Conn {
38 - return s.conn
39 -}
40 -
41 -// wire pipe between two network conns. yay io.
42 -func newStreamPair(n1 *Network, p2 peer.Peer) (*Stream, *Stream) {
43 - p1 := n1.local
44 - r1, w1 := io.Pipe()
45 - r2, w2 := io.Pipe()
46 -
47 - s1 := &Stream{Reader: r1, Writer: w2}
48 - s2 := &Stream{Reader: r2, Writer: w1}
49 -
50 - n1.Lock()
51 - n1.conns[p2].addStream(s1)
52 - n2 := n1.conns[p2].remote
53 - n1.Unlock()
54 -
55 - n2.Lock()
56 - n2.conns[p1].addStream(s2)
57 - n2.Unlock()
58 - n2.handle(s2)
59 -
60 - return s1, s2
61 -}
62 -
63 -type Conn struct {
64 - connected bool
65 - local *Network
66 - remote *Network
67 - streams []*Stream
68 - sync.RWMutex
69 -}
70 -
71 -func (c *Conn) Close() error {
72 - c.Lock()
73 - defer c.Unlock()
74 -
75 - c.connected = false
76 - for _, s := range c.streams {
77 - go s.Close()
78 - }
79 - c.streams = nil
80 - return nil
81 -}
82 -
83 -func (c *Conn) addStream(s *Stream) {
84 - c.Lock()
85 - defer c.Unlock()
86 -
87 - s.conn = c
88 - c.streams = append(c.streams, s)
89 -}
90 -
91 -func (c *Conn) removeStream(s *Stream) {
92 - c.Lock()
93 - defer c.Unlock()
94 -
95 - strs := make([]*Stream, 0, len(c.streams))
96 - for _, s2 := range c.streams {
97 - if s2 != s {
98 - strs = append(strs, s2)
99 - }
100 - }
101 -}
102 -
103 -func (c *Conn) NewStreamWithProtocol(pr inet.ProtocolID, p peer.Peer) (inet.Stream, error) {
104 -
105 - if _, connected := c.local.conns[p]; !connected {
106 - return nil, fmt.Errorf("cannot create new stream for %s. not connected.", p)
107 - }
108 -
109 - log.Debugf("NewStreamWithProtocol: %s --> %s", c.local, p)
110 - ss, _ := newStreamPair(c.local, p)
111 -
112 - if err := inet.WriteProtocolHeader(pr, ss); err != nil {
113 - ss.Close()
114 - return nil, err
115 - }
116 -
117 - return ss, nil
118 -}
119 -
120 -// LocalMultiaddr is the Multiaddr on this side
121 -func (c *Conn) LocalMultiaddr() ma.Multiaddr {
122 - return nil
123 -}
124 -
125 -// LocalPeer is the Peer on our side of the connection
126 -func (c *Conn) LocalPeer() peer.Peer {
127 - return c.local.local
128 -}
129 -
130 -// RemoteMultiaddr is the Multiaddr on the remote side
131 -func (c *Conn) RemoteMultiaddr() ma.Multiaddr {
132 - return nil
133 -}
134 -
135 -// RemotePeer is the Peer on the remote side
136 -func (c *Conn) RemotePeer() peer.Peer {
137 - return c.remote.local
138 -}
139 -
140 -// network implements the Network interface,
141 -type Network struct {
142 - local peer.Peer // local peer
143 - mux inet.Mux // protocol multiplexing
144 -
145 - conns map[peer.Peer]*Conn
146 - sync.RWMutex
147 -
148 - cg ctxgroup.ContextGroup // for Context closing
149 -}
150 -
151 -func MakeNetworks(ctx context.Context, peers []peer.Peer) (nets []*Network, err error) {
152 - nets = make([]*Network, len(peers))
153 - for i, p := range peers {
154 - ps := peer.NewPeerstore()
155 - nets[i], err = newNetwork(ctx, p, ps)
156 - if err != nil {
157 - return nil, err
158 - }
159 - }
160 -
161 - i := 0
162 - for _, n1 := range nets {
163 - for _, n2 := range nets {
164 - n1.conns[n2.local] = &Conn{local: n1, remote: n2}
165 - log.Debugf("%d setup %s -> %s", i, n1, n2)
166 - i++
167 - }
168 - }
169 -
170 - return nets, nil
171 -}
172 -
173 -// NewNetwork constructs a new Mock network
174 -func newNetwork(ctx context.Context, local peer.Peer, peers peer.Peerstore) (*Network, error) {
175 -
176 - n := &Network{
177 - local: local,
178 - mux: inet.Mux{Handlers: inet.StreamHandlerMap{}},
179 - cg: ctxgroup.WithContext(ctx),
180 - conns: map[peer.Peer]*Conn{},
181 - }
182 -
183 - n.cg.SetTeardown(n.close)
184 - return n, nil
185 -}
186 -func (n *Network) String() string {
187 - return fmt.Sprintf("<Network %s - %d conns>", n.local, len(n.conns))
188 -}
189 -
190 -func (n *Network) handle(s inet.Stream) {
191 - go n.mux.Handle(s)
192 -}
193 -
194 -// DialPeer attempts to establish a connection to a given peer.
195 -// Respects the context.
196 -func (n *Network) DialPeer(ctx context.Context, p peer.Peer) error {
197 - n.Lock()
198 - defer n.Unlock()
199 -
200 - c, ok := n.conns[p]
201 - if !ok {
202 - return fmt.Errorf("cannot connect to %s (mock needs all nets at start)", p)
203 - }
204 - c.connected = true
205 - return nil
206 -}
207 -
208 -// CtxGroup returns the network's ContextGroup
209 -func (n *Network) CtxGroup() ctxgroup.ContextGroup {
210 - return n.cg
211 -}
212 -
213 -// LocalPeer the network's LocalPeer
214 -func (n *Network) LocalPeer() peer.Peer {
215 - return n.local
216 -}
217 -
218 -// Peers returns the connected peers
219 -func (n *Network) Peers() []peer.Peer {
220 - n.RLock()
221 - defer n.RUnlock()
222 -
223 - peers := make([]peer.Peer, 0, len(n.conns))
224 - for _, c := range n.conns {
225 - if c.connected {
226 - peers = append(peers, c.RemotePeer())
227 - }
228 - }
229 - return peers
230 -}
231 -
232 -// Conns returns the connected peers
233 -func (n *Network) Conns() []inet.Conn {
234 - n.RLock()
235 - defer n.RUnlock()
236 -
237 - out := make([]inet.Conn, 0, len(n.conns))
238 - for _, c := range n.conns {
239 - if c.connected {
240 - out = append(out, c)
241 - }
242 - }
243 - return out
244 -}
245 -
246 -// ClosePeer connection to peer
247 -func (n *Network) ClosePeer(p peer.Peer) error {
248 - c, ok := n.conns[p]
249 - if !ok {
250 - return nil
251 - }
252 - return c.Close()
253 -}
254 -
255 -// close is the real teardown function
256 -func (n *Network) close() error {
257 - for _, c := range n.conns {
258 - c.Close()
259 - }
260 - return nil
261 -}
262 -
263 -// Close calls the ContextCloser func
264 -func (n *Network) Close() error {
265 - return n.cg.Close()
266 -}
267 -
268 -// BandwidthTotals returns the total amount of bandwidth transferred
269 -func (n *Network) BandwidthTotals() (in uint64, out uint64) {
270 - // need to implement this. probably best to do it in swarm this time.
271 - // need a "metrics" object
272 - return 0, 0
273 -}
274 -
275 -// ListenAddresses returns a list of addresses at which this network listens.
276 -func (n *Network) ListenAddresses() []ma.Multiaddr {
277 - return []ma.Multiaddr{}
278 -}
279 -
280 -// InterfaceListenAddresses returns a list of addresses at which this network
281 -// listens. It expands "any interface" addresses (/ip4/0.0.0.0, /ip6/::) to
282 -// use the known local interfaces.
283 -func (n *Network) InterfaceListenAddresses() ([]ma.Multiaddr, error) {
284 - return []ma.Multiaddr{}, nil
285 -}
286 -
287 -// Connectedness returns a state signaling connection capabilities
288 -// For now only returns Connecter || NotConnected. Expand into more later.
289 -func (n *Network) Connectedness(p peer.Peer) inet.Connectedness {
290 - n.Lock()
291 - defer n.Unlock()
292 -
293 - if _, found := n.conns[p]; found && n.conns[p].connected {
294 - return inet.Connected
295 - }
296 - return inet.NotConnected
297 -}
298 -
299 -// NewStream returns a new stream to given peer p.
300 -// If there is no connection to p, attempts to create one.
301 -// If ProtocolID is "", writes no header.
302 -func (c *Network) NewStream(pr inet.ProtocolID, p peer.Peer) (inet.Stream, error) {
303 - return c.conns[p].NewStreamWithProtocol(pr, p)
304 -}
305 -
306 -// SetHandler sets the protocol handler on the Network's Muxer.
307 -// This operation is threadsafe.
308 -func (n *Network) SetHandler(p inet.ProtocolID, h inet.StreamHandler) {
309 - n.mux.SetHandler(p, h)
310 -}
net/mock/mock_test.go deleted
-194
@@ -1,194 +0,0 @@
1 -package mocknet
2 -
3 -import (
4 - "bytes"
5 - "io"
6 - "math/rand"
7 - "sync"
8 - "testing"
9 -
10 - inet "github.com/jbenet/go-ipfs/net"
11 - peer "github.com/jbenet/go-ipfs/peer"
12 -
13 - context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14 - testutil "github.com/jbenet/go-ipfs/util/testutil"
15 -)
16 -
17 -func TestNetworkSetup(t *testing.T) {
18 -
19 - p1 := testutil.RandPeer()
20 - p2 := testutil.RandPeer()
21 - p3 := testutil.RandPeer()
22 - peers := []peer.Peer{p1, p2, p3}
23 -
24 - nets, err := MakeNetworks(context.Background(), peers)
25 - if err != nil {
26 - t.Fatal(err)
27 - }
28 -
29 - // check things
30 -
31 - if len(nets) != 3 {
32 - t.Error("nets must be 3")
33 - }
34 -
35 - for i, n := range nets {
36 - if n.local != peers[i] {
37 - t.Error("peer mismatch")
38 - }
39 -
40 - if len(n.conns) != len(nets) {
41 - t.Error("conn mismatch")
42 - }
43 -
44 - for _, c := range n.conns {
45 - if c.remote.conns[n.local] == nil {
46 - t.Error("conn other side fail")
47 - }
48 - if c.remote.conns[n.local].remote.local != n.local {
49 - t.Error("conn other side fail")
50 - }
51 - }
52 -
53 - }
54 -
55 -}
56 -
57 -func TestStreams(t *testing.T) {
58 -
59 - p1 := testutil.RandPeer()
60 - p2 := testutil.RandPeer()
61 - p3 := testutil.RandPeer()
62 - peers := []peer.Peer{p1, p2, p3}
63 -
64 - nets, err := MakeNetworks(context.Background(), peers)
65 - if err != nil {
66 - t.Fatal(err)
67 - }
68 -
69 - nets[1].SetHandler(inet.ProtocolDHT, func(s inet.Stream) {
70 - go func() {
71 - b := make([]byte, 4)
72 - if _, err := io.ReadFull(s, b); err != nil {
73 - panic(err)
74 - }
75 - if !bytes.Equal(b, []byte("beep")) {
76 - panic("bytes mismatch")
77 - }
78 - if _, err := s.Write([]byte("boop")); err != nil {
79 - panic(err)
80 - }
81 - s.Close()
82 - }()
83 - })
84 -
85 - s, err := nets[0].NewStream(inet.ProtocolDHT, nets[1].local)
86 - if err != nil {
87 - t.Fatal(err)
88 - }
89 -
90 - if _, err := s.Write([]byte("beep")); err != nil {
91 - panic(err)
92 - }
93 - b := make([]byte, 4)
94 - if _, err := io.ReadFull(s, b); err != nil {
95 - panic(err)
96 - }
97 - if !bytes.Equal(b, []byte("boop")) {
98 - panic("bytes mismatch 2")
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 -}