updated Conn and Swarm
This Commit changes the relationship between Conn and Swarm. After this, Conn is significantly more autonomous, and follows an interface. From here, it will be very easy to make the MultiConn (that handles multiple Conns per peer).
Juan Batiz-Benet committed
Oct 16, 2014 at 07:37 UTC
1edc5a4613d71b6aa02fb6c61b5005c3dbd4dd31
4 files changed
+53
-103
net/conn/conn.go
+2
-1
@@ -52,7 +52,7 @@ type singleConn struct {
52
}
53
54
// Map maps Keys (Peer.IDs) to Connections.
55
-type Map map[u.Key]*Conn
55
+type Map map[u.Key]Conn
56
57
// newConn constructs a new connection
58
func newSingleConn(ctx context.Context, local, remote *peer.Peer,
@@ -171,6 +171,7 @@ func (c *singleConn) waitToClose(ctx context.Context) {
171
// closing channels
172
c.insecure.outgoing.Close()
173
c.secure.Close()
174
+ close(c.msgpipe.Incoming)
175
}
176
177
// IsOpen returns whether this Conn is open or closed.
net/swarm/conn.go
+39
-94
@@ -4,14 +4,12 @@ import (
4
"errors"
5
"fmt"
6
7
- spipe "github.com/jbenet/go-ipfs/crypto/spipe"
7
conn "github.com/jbenet/go-ipfs/net/conn"
8
handshake "github.com/jbenet/go-ipfs/net/handshake"
9
msg "github.com/jbenet/go-ipfs/net/message"
10
11
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
12
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
14
- manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
13
)
14
15
// Open listeners for each network the swarm should listen on
@@ -39,7 +37,8 @@ func (s *Swarm) listen() error {
37
38
// Listen for new connections on the given multiaddr
39
func (s *Swarm) connListen(maddr ma.Multiaddr) error {
42
- list, err := manet.Listen(maddr)
40
+
41
+ list, err := conn.Listen(s.ctx, maddr, s.local, s.peers)
42
if err != nil {
43
return err
44
}
@@ -55,17 +54,12 @@ func (s *Swarm) connListen(maddr ma.Multiaddr) error {
54
// Accept and handle new connections on this listener until it errors
55
go func() {
56
for {
58
- nconn, err := list.Accept()
59
- if err != nil {
60
- e := fmt.Errorf("Failed to accept connection: %s - %s", maddr, err)
61
- s.errChan <- e
57
+ select {
58
+ case <-s.ctx.Done():
59
+ return
60
63
- // if cancel is nil, we're closed.
64
- if s.cancel == nil {
65
- return
66
- }
67
- } else {
68
- go s.handleIncomingConn(nconn)
61
+ case conn := <-list.Accept():
62
+ go s.handleIncomingConn(conn)
63
}
64
}
65
}()
@@ -74,42 +68,24 @@ func (s *Swarm) connListen(maddr ma.Multiaddr) error {
68
}
69
70
// Handle getting ID from this peer, handshake, and adding it into the map
77
-func (s *Swarm) handleIncomingConn(nconn manet.Conn) {
78
-
79
- // Construct conn with nil peer for now, because we don't know its ID yet.
80
- // connSetup will figure this out, and pull out / construct the peer.
81
- c, err := conn.NewConn(s.local, nil, nconn)
82
- if err != nil {
83
- s.errChan <- err
84
- return
85
- }
71
+func (s *Swarm) handleIncomingConn(nconn conn.Conn) {
72
73
// Setup the new connection
88
- err = s.connSetup(c)
74
+ err := s.connSetup(nconn)
75
if err != nil && err != ErrAlreadyOpen {
76
s.errChan <- err
91
- c.Close()
77
+ nconn.Close()
78
}
79
}
80
81
// connSetup adds the passed in connection to its peerMap and starts
82
// the fanIn routine for that connection
97
-func (s *Swarm) connSetup(c *conn.Conn) error {
83
+func (s *Swarm) connSetup(c conn.Conn) error {
84
if c == nil {
85
return errors.New("Tried to start nil connection.")
86
}
87
102
- if c.Remote != nil {
103
- log.Debug("%s Starting connection: %s", c.Local, c.Remote)
104
- } else {
105
- log.Debug("%s Starting connection: [unknown peer]", c.Local)
106
- }
107
-
108
- if err := s.connSecure(c); err != nil {
109
- return fmt.Errorf("Conn securing error: %v", err)
110
- }
111
-
112
- log.Debug("%s secured connection: %s", c.Local, c.Remote)
88
+ log.Debug("%s Started connection: %s", c.LocalPeer(), c.RemotePeer())
89
90
// add address of connection to Peer. Maybe it should happen in connSecure.
91
// NOT adding this address here, because the incoming address in TCP
@@ -123,12 +99,12 @@ func (s *Swarm) connSetup(c *conn.Conn) error {
99
100
// add to conns
101
s.connsLock.Lock()
126
- if _, ok := s.conns[c.Remote.Key()]; ok {
102
+ if _, ok := s.conns[c.RemotePeer().Key()]; ok {
103
log.Debug("Conn already open!")
104
s.connsLock.Unlock()
105
return ErrAlreadyOpen
106
}
131
- s.conns[c.Remote.Key()] = c
107
+ s.conns[c.RemotePeer().Key()] = c
108
log.Debug("Added conn to map!")
109
s.connsLock.Unlock()
110
@@ -137,77 +113,51 @@ func (s *Swarm) connSetup(c *conn.Conn) error {
113
return nil
114
}
115
140
-// connSecure setups a secure remote connection.
141
-func (s *Swarm) connSecure(c *conn.Conn) error {
142
-
143
- sp, err := spipe.NewSecurePipe(s.ctx, 10, s.local, s.peers)
144
- if err != nil {
145
- return err
146
- }
147
-
148
- err = sp.Wrap(s.ctx, spipe.Duplex{
149
- In: c.Incoming.MsgChan,
150
- Out: c.Outgoing.MsgChan,
151
- })
152
- if err != nil {
153
- return err
154
- }
155
-
156
- if c.Remote == nil {
157
- c.Remote = sp.RemotePeer()
158
-
159
- } else if c.Remote != sp.RemotePeer() {
160
- panic("peers not being constructed correctly.")
161
- }
162
-
163
- c.Secure = sp
164
- return nil
165
-}
166
-
116
// connVersionExchange exchanges local and remote versions and compares them
117
// closes remote and returns an error in case of major difference
169
-func (s *Swarm) connVersionExchange(remote *conn.Conn) error {
170
- var remoteHandshake, localHandshake *handshake.Handshake1
171
- localHandshake = handshake.CurrentHandshake()
118
+func (s *Swarm) connVersionExchange(r conn.Conn) error {
119
+ rpeer := r.RemotePeer()
120
173
- myVerBytes, err := proto.Marshal(localHandshake)
121
+ var remoteH, localH *handshake.Handshake1
122
+ localH = handshake.CurrentHandshake()
123
+
124
+ myVerBytes, err := proto.Marshal(localH)
125
if err != nil {
126
return err
127
}
128
178
- remote.Secure.Out <- myVerBytes
179
-
180
- log.Debug("Send my version(%s) [to = %s]", localHandshake, remote.Peer)
129
+ r.MsgOut() <- msg.New(rpeer, myVerBytes)
130
+ log.Debug("Sent my version(%s) [to = %s]", localH, rpeer)
131
132
select {
133
case <-s.ctx.Done():
134
return s.ctx.Err()
135
186
- case <-remote.Closed:
187
- return errors.New("remote closed connection during version exchange")
136
+ // case <-remote.Done():
137
+ // return errors.New("remote closed connection during version exchange")
138
189
- case data, ok := <-remote.Secure.In:
139
+ case data, ok := <-r.MsgIn():
140
if !ok {
191
- return fmt.Errorf("Error retrieving from conn: %v", remote.Peer)
141
+ return fmt.Errorf("Error retrieving from conn: %v", rpeer)
142
}
143
194
- remoteHandshake = new(handshake.Handshake1)
195
- err = proto.Unmarshal(data, remoteHandshake)
144
+ remoteH = new(handshake.Handshake1)
145
+ err = proto.Unmarshal(data.Data(), remoteH)
146
if err != nil {
147
s.Close()
148
return fmt.Errorf("connSetup: could not decode remote version: %q", err)
149
}
150
201
- log.Debug("Received remote version(%s) [from = %s]", remoteHandshake, remote.Peer)
151
+ log.Debug("Received remote version(%s) [from = %s]", remoteH, rpeer)
152
}
153
204
- if err := handshake.Compatible(localHandshake, remoteHandshake); err != nil {
205
- log.Info("%s (%s) incompatible version with %s (%s)", s.local, localHandshake, remote.Peer, remoteHandshake)
206
- remote.Close()
154
+ if err := handshake.Compatible(localH, remoteH); err != nil {
155
+ log.Info("%s (%s) incompatible version with %s (%s)", s.local, localH, rpeer, remoteH)
156
+ r.Close()
157
return err
158
}
159
210
- log.Debug("[peer: %s] Version compatible", remote.Peer)
160
+ log.Debug("[peer: %s] Version compatible", rpeer)
161
return nil
162
}
163
@@ -237,14 +187,14 @@ func (s *Swarm) fanOut() {
187
// log.Debug("[peer: %s] Sent message [to = %s]", s.local, msg.Peer())
188
189
// queue it in the connection's buffer
240
- conn.Secure.Out <- msg.Data()
190
+ conn.MsgOut() <- msg
191
}
192
}
193
}
194
195
// Handles the receiving + wrapping of messages, per conn.
196
// Consider using reflect.Select with one goroutine instead of n.
247
-func (s *Swarm) fanIn(c *conn.Conn) {
197
+func (s *Swarm) fanIn(c conn.Conn) {
198
for {
199
select {
200
case <-s.ctx.Done():
@@ -252,26 +202,21 @@ func (s *Swarm) fanIn(c *conn.Conn) {
202
c.Close()
203
goto out
204
255
- case <-c.Closed:
256
- goto out
257
-
258
- case data, ok := <-c.Secure.In:
205
+ case data, ok := <-c.MsgIn():
206
if !ok {
260
- e := fmt.Errorf("Error retrieving from conn: %v", c.Remote)
207
+ e := fmt.Errorf("Error retrieving from conn: %v", c.RemotePeer())
208
s.errChan <- e
209
goto out
210
}
211
212
// log.Debug("[peer: %s] Received message [from = %s]", s.local, c.Peer)
266
-
267
- msg := msg.New(c.Remote, data)
268
- s.Incoming <- msg
213
+ s.Incoming <- data
214
}
215
}
216
217
out:
218
s.connsLock.Lock()
274
- delete(s.conns, c.Remote.Key())
219
+ delete(s.conns, c.RemotePeer().Key())
220
s.connsLock.Unlock()
221
}
222
net/swarm/swarm.go
+10
-6
@@ -11,7 +11,6 @@ import (
11
u "github.com/jbenet/go-ipfs/util"
12
13
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
14
- manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
14
)
15
16
var log = u.Logger("swarm")
@@ -61,7 +60,7 @@ type Swarm struct {
60
connsLock sync.RWMutex
61
62
// listeners for each network address
64
- listeners []manet.Listener
63
+ listeners []conn.Listener
64
65
// cancel is an internal function used to stop the Swarm's processing.
66
cancel context.CancelFunc
@@ -110,7 +109,7 @@ func (s *Swarm) Close() error {
109
// etc. to achive connection.
110
//
111
// For now, Dial uses only TCP. This will be extended.
113
-func (s *Swarm) Dial(peer *peer.Peer) (*conn.Conn, error) {
112
+func (s *Swarm) Dial(peer *peer.Peer) (conn.Conn, error) {
113
if peer.ID.Equal(s.local.ID) {
114
return nil, errors.New("Attempted connection to self!")
115
}
@@ -128,7 +127,12 @@ func (s *Swarm) Dial(peer *peer.Peer) (*conn.Conn, error) {
127
}
128
129
// open connection to peer
131
- c, err = conn.Dial("tcp", s.local, peer)
130
+ d := &conn.Dialer{
131
+ LocalPeer: s.local,
132
+ Peerstore: s.peers,
133
+ }
134
+
135
+ c, err = d.Dial(s.ctx, "tcp", s.local)
136
if err != nil {
137
return nil, err
138
}
@@ -142,7 +146,7 @@ func (s *Swarm) Dial(peer *peer.Peer) (*conn.Conn, error) {
146
}
147
148
// GetConnection returns the connection in the swarm to given peer.ID
145
-func (s *Swarm) GetConnection(pid peer.ID) *conn.Conn {
149
+func (s *Swarm) GetConnection(pid peer.ID) conn.Conn {
150
s.connsLock.RLock()
151
c, found := s.conns[u.Key(pid)]
152
s.connsLock.RUnlock()
@@ -181,7 +185,7 @@ func (s *Swarm) GetPeerList() []*peer.Peer {
185
var out []*peer.Peer
186
s.connsLock.RLock()
187
for _, p := range s.conns {
184
- out = append(out, p.Remote)
188
+ out = append(out, p.RemotePeer())
189
}
190
s.connsLock.RUnlock()
191
return out
net/swarm/swarm_test.go
+2
-2
@@ -49,13 +49,13 @@ func setupPeer(t *testing.T, addr string) *peer.Peer {
49
p.PrivKey = sk
50
p.PubKey = pk
51
p.AddAddress(tcp)
52
- return p, nil
52
+ return p
53
}
54
55
func makeSwarms(ctx context.Context, t *testing.T, peers map[string]string) []*Swarm {
56
swarms := []*Swarm{}
57
58
- for key, addr := range peers {
58
+ for _, addr := range peers {
59
local := setupPeer(t, addr)
60
peerstore := peer.NewPeerstore()
61
swarm, err := NewSwarm(ctx, local, peerstore)