@cryptotaxi247 / kubo / commits / f8d70f344

simultaneous open should work for now

It's a patch, really. it's not the full multiconn fix.

Juan Batiz-Benet committed Oct 18, 2014 at 19:57 UTC f8d70f344b2366db8c3376ee4429a5e57020f63e
5 files changed +134 -32
crypto/spipe/handshake.go
+1 -1
@@ -205,7 +205,7 @@ func (s *SecurePipe) handshake() error {
205 }
206
207 if bytes.Compare(resp2, finished) != 0 {
208 - return errors.New("Negotiation failed.")
208 + return fmt.Errorf("Negotiation failed, got: %s", resp2)
209 }
210
211 log.Debug("%s handshake: Got node id: %s", s.local, s.remote)
net/swarm/conn.go
+9 -6
@@ -69,7 +69,7 @@ func (s *Swarm) connListen(maddr ma.Multiaddr) error {
69 func (s *Swarm) handleIncomingConn(nconn conn.Conn) {
70
71 // Setup the new connection
72 - err := s.connSetup(nconn)
72 + _, err := s.connSetup(nconn)
73 if err != nil && err != ErrAlreadyOpen {
74 s.errChan <- err
75 nconn.Close()
@@ -78,9 +78,9 @@ func (s *Swarm) handleIncomingConn(nconn conn.Conn) {
78
79 // connSetup adds the passed in connection to its peerMap and starts
80 // the fanIn routine for that connection
81 -func (s *Swarm) connSetup(c conn.Conn) error {
81 +func (s *Swarm) connSetup(c conn.Conn) (conn.Conn, error) {
82 if c == nil {
83 - return errors.New("Tried to start nil connection.")
83 + return nil, errors.New("Tried to start nil connection.")
84 }
85
86 log.Debug("%s Started connection: %s", c.LocalPeer(), c.RemotePeer())
@@ -93,10 +93,13 @@ func (s *Swarm) connSetup(c conn.Conn) error {
93
94 // add to conns
95 s.connsLock.Lock()
96 - if _, ok := s.conns[c.RemotePeer().Key()]; ok {
96 + if c2, ok := s.conns[c.RemotePeer().Key()]; ok {
97 log.Debug("Conn already open!")
98 s.connsLock.Unlock()
99 - return ErrAlreadyOpen
99 +
100 + c.Close()
101 + return c2, nil // not error anymore, use existing conn.
102 + // return ErrAlreadyOpen
103 }
104 s.conns[c.RemotePeer().Key()] = c
105 log.Debug("Added conn to map!")
@@ -104,7 +107,7 @@ func (s *Swarm) connSetup(c conn.Conn) error {
107
108 // kick off reader goroutine
109 go s.fanIn(c)
107 - return nil
110 + return c, nil
111 }
112
113 // Handles the unwrapping + sending of messages to the right connection.
net/swarm/simul_test.go new
+73
@@ -0,0 +1,73 @@
1 +package swarm
2 +
3 +import (
4 + "fmt"
5 + "sync"
6 + "testing"
7 +
8 + peer "github.com/jbenet/go-ipfs/peer"
9 +
10 + context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11 +)
12 +
13 +func TestSimultOpen(t *testing.T) {
14 + t.Skip("skipping for another test")
15 +
16 + addrs := []string{
17 + "/ip4/127.0.0.1/tcp/1244",
18 + "/ip4/127.0.0.1/tcp/1245",
19 + }
20 +
21 + ctx := context.Background()
22 + swarms, _ := makeSwarms(ctx, t, addrs)
23 +
24 + // connect everyone
25 + {
26 + var wg sync.WaitGroup
27 + connect := func(s *Swarm, dst *peer.Peer) {
28 + // copy for other peer
29 + cp := &peer.Peer{ID: dst.ID}
30 + cp.AddAddress(dst.Addresses[0])
31 +
32 + if _, err := s.Dial(cp); err != nil {
33 + t.Fatal("error swarm dialing to peer", err)
34 + }
35 + log.Info("done?!?")
36 + wg.Done()
37 + }
38 +
39 + log.Info("Connecting swarms simultaneously.")
40 + wg.Add(2)
41 + go connect(swarms[0], swarms[1].local)
42 + go connect(swarms[1], swarms[0].local)
43 + wg.Wait()
44 + }
45 +
46 + for _, s := range swarms {
47 + s.Close()
48 + }
49 +}
50 +
51 +func TestSimultOpenMany(t *testing.T) {
52 + t.Skip("laggy")
53 +
54 + addrs := []string{}
55 + for i := 2200; i < 2300; i++ {
56 + s := fmt.Sprintf("/ip4/127.0.0.1/tcp/%d", i)
57 + addrs = append(addrs, s)
58 + }
59 +
60 + SubtestSwarm(t, addrs, 10)
61 +}
62 +
63 +func TestSimultOpenFewStress(t *testing.T) {
64 +
65 + for i := 0; i < 100; i++ {
66 + addrs := []string{
67 + fmt.Sprintf("/ip4/127.0.0.1/tcp/%d", 1900+i),
68 + fmt.Sprintf("/ip4/127.0.0.1/tcp/%d", 2900+i),
69 + }
70 +
71 + SubtestSwarm(t, addrs, 10)
72 + }
73 +}
net/swarm/swarm.go
+2 -1
@@ -137,7 +137,8 @@ func (s *Swarm) Dial(peer *peer.Peer) (conn.Conn, error) {
137 return nil, err
138 }
139
140 - if err := s.connSetup(c); err != nil {
140 + c, err = s.connSetup(c)
141 + if err != nil {
142 c.Close()
143 return nil, err
144 }
net/swarm/swarm_test.go
+49 -24
@@ -2,7 +2,7 @@ package swarm
2
3 import (
4 "bytes"
5 - "fmt"
5 + "sync"
6 "testing"
7 "time"
8
@@ -23,6 +23,7 @@ func pong(ctx context.Context, swarm *Swarm) {
23 case m1 := <-swarm.Incoming:
24 if bytes.Equal(m1.Data(), []byte("ping")) {
25 m2 := msg.New(m1.Peer(), []byte("pong"))
26 + log.Debug("%s pong %s", swarm.local, m1.Peer())
27 swarm.Outgoing <- m2
28 }
29 }
@@ -52,10 +53,10 @@ func setupPeer(t *testing.T, addr string) *peer.Peer {
53 return p
54 }
55
55 -func makeSwarms(ctx context.Context, t *testing.T, peers map[string]string) []*Swarm {
56 +func makeSwarms(ctx context.Context, t *testing.T, addrs []string) ([]*Swarm, []*peer.Peer) {
57 swarms := []*Swarm{}
58
58 - for _, addr := range peers {
59 + for _, addr := range addrs {
60 local := setupPeer(t, addr)
61 peerstore := peer.NewPeerstore()
62 swarm, err := NewSwarm(ctx, local, peerstore)
@@ -65,35 +66,46 @@ func makeSwarms(ctx context.Context, t *testing.T, peers map[string]string) []*S
66 swarms = append(swarms, swarm)
67 }
68
68 - return swarms
69 + peers := make([]*peer.Peer, len(swarms))
70 + for i, s := range swarms {
71 + peers[i] = s.local
72 + }
73 +
74 + return swarms, peers
75 }
76
71 -func TestSwarm(t *testing.T) {
72 - peers := map[string]string{
73 - "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a30": "/ip4/127.0.0.1/tcp/1234",
74 - "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31": "/ip4/127.0.0.1/tcp/2345",
75 - "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a32": "/ip4/127.0.0.1/tcp/3456",
76 - "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33": "/ip4/127.0.0.1/tcp/4567",
77 - "11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a34": "/ip4/127.0.0.1/tcp/5678",
78 - }
77 +func SubtestSwarm(t *testing.T, addrs []string, MsgNum int) {
78 + // t.Skip("skipping for another test")
79
80 ctx := context.Background()
81 - swarms := makeSwarms(ctx, t, peers)
81 + swarms, peers := makeSwarms(ctx, t, addrs)
82
83 // connect everyone
84 - for _, s := range swarms {
85 - peers, err := s.peers.All()
86 - if err != nil {
87 - t.Fatal(err)
84 + {
85 + var wg sync.WaitGroup
86 + connect := func(s *Swarm, dst *peer.Peer) {
87 + // copy for other peer
88 + cp := &peer.Peer{ID: dst.ID}
89 + cp.AddAddress(dst.Addresses[0])
90 +
91 + log.Info("SWARM TEST: %s dialing %s", s.local, dst)
92 + if _, err := s.Dial(cp); err != nil {
93 + t.Fatal("error swarm dialing to peer", err)
94 + }
95 + log.Info("SWARM TEST: %s connected to %s", s.local, dst)
96 + wg.Done()
97 }
98
90 - for _, p := range *peers {
91 - fmt.Println("dialing")
92 - if _, err := s.Dial(p); err != nil {
93 - t.Fatal("error swarm dialing to peer", err)
99 + log.Info("Connecting swarms simultaneously.")
100 + for _, s := range swarms {
101 + for _, p := range peers {
102 + if p != s.local { // don't connect to self.
103 + wg.Add(1)
104 + connect(s, p)
105 + }
106 }
95 - fmt.Println("dialed")
107 }
108 + wg.Wait()
109 }
110
111 // ping/pong
@@ -114,9 +126,9 @@ func TestSwarm(t *testing.T) {
126 t.Fatal(err)
127 }
128
117 - MsgNum := 1000
129 for k := 0; k < MsgNum; k++ {
130 for _, p := range *peers {
131 + log.Debug("%s ping %s", s1.local, p)
132 s1.Outgoing <- msg.New(p, []byte("ping"))
133 }
134 }
@@ -143,10 +155,23 @@ func TestSwarm(t *testing.T) {
155 }
156
157 cancel()
146 - <-time.After(50 * time.Millisecond)
158 + <-time.After(50 * time.Microsecond)
159 }
160
161 for _, s := range swarms {
162 s.Close()
163 }
164 }
165 +
166 +func TestSwarm(t *testing.T) {
167 +
168 + addrs := []string{
169 + "/ip4/127.0.0.1/tcp/1234",
170 + "/ip4/127.0.0.1/tcp/1235",
171 + "/ip4/127.0.0.1/tcp/1236",
172 + "/ip4/127.0.0.1/tcp/1237",
173 + "/ip4/127.0.0.1/tcp/1238",
174 + }
175 +
176 + SubtestSwarm(t, addrs, 1000)
177 +}