@cryptotaxi247 / kubo / commits / 1b1ef6aa0

add local to net/conn

Juan Batiz-Benet committed Oct 11, 2014 at 04:31 UTC 1b1ef6aa090472434d977984f02c53be5c236801
4 files changed +79 -36
net/conn/conn.go
+24 -15
@@ -4,7 +4,6 @@ import (
4 "fmt"
5
6 msgio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-msgio"
7 - ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
7 manet "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr/net"
8
9 spipe "github.com/jbenet/go-ipfs/crypto/spipe"
@@ -22,9 +21,9 @@ const MaxMessageSize = 1 << 20
21
22 // Conn represents a connection to another Peer (IPFS Node).
23 type Conn struct {
25 - Peer *peer.Peer
26 - Addr ma.Multiaddr
27 - Conn manet.Conn
24 + Local *peer.Peer
25 + Remote *peer.Peer
26 + Conn manet.Conn
27
28 Closed chan bool
29 Outgoing *msgio.Chan
@@ -36,11 +35,11 @@ type Conn struct {
35 type Map map[u.Key]*Conn
36
37 // NewConn constructs a new connection
39 -func NewConn(peer *peer.Peer, addr ma.Multiaddr, mconn manet.Conn) (*Conn, error) {
38 +func NewConn(local, remote *peer.Peer, mconn manet.Conn) (*Conn, error) {
39 conn := &Conn{
41 - Peer: peer,
42 - Addr: addr,
43 - Conn: mconn,
40 + Local: local,
41 + Remote: remote,
42 + Conn: mconn,
43 }
44
45 if err := conn.newChans(); err != nil {
@@ -52,18 +51,28 @@ func NewConn(peer *peer.Peer, addr ma.Multiaddr, mconn manet.Conn) (*Conn, error
51
52 // Dial connects to a particular peer, over a given network
53 // Example: Dial("udp", peer)
55 -func Dial(network string, peer *peer.Peer) (*Conn, error) {
56 - addr := peer.NetAddress(network)
57 - if addr == nil {
58 - return nil, fmt.Errorf("No address for network %s", network)
54 +func Dial(network string, local, remote *peer.Peer) (*Conn, error) {
55 + laddr := local.NetAddress(network)
56 + if laddr == nil {
57 + return nil, fmt.Errorf("No local address for network %s", network)
58 }
59
61 - nconn, err := manet.Dial(addr)
60 + raddr := remote.NetAddress(network)
61 + if raddr == nil {
62 + return nil, fmt.Errorf("No remote address for network %s", network)
63 + }
64 +
65 + // TODO: try to get reusing addr/ports to work.
66 + // dialer := manet.Dialer{LocalAddr: laddr}
67 + dialer := manet.Dialer{}
68 +
69 + log.Info("%s %s dialing %s %s", local, laddr, remote, raddr)
70 + nconn, err := dialer.Dial(raddr)
71 if err != nil {
72 return nil, err
73 }
74
66 - return NewConn(peer, addr, nconn)
75 + return NewConn(local, remote, nconn)
76 }
77
78 // Construct new channels for given Conn.
@@ -84,7 +93,7 @@ func (c *Conn) newChans() error {
93
94 // Close closes the connection, and associated channels.
95 func (c *Conn) Close() error {
87 - log.Debug("Closing Conn with %v", c.Peer)
96 + log.Debug("%s closing Conn with %s", c.Local, c.Remote)
97 if c.Conn == nil {
98 return fmt.Errorf("Already closed") // already closed
99 }
net/conn/conn_test.go
+7 -2
@@ -65,12 +65,17 @@ func TestDial(t *testing.T) {
65 }
66 go echoListen(listener)
67
68 - p, err := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33", "/ip4/127.0.0.1/tcp/1234")
68 + p1, err := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33", "/ip4/127.0.0.1/tcp/1234")
69 if err != nil {
70 t.Fatal("error setting up peer", err)
71 }
72
73 - c, err := Dial("tcp", p)
73 + p2, err := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a34", "/ip4/127.0.0.1/tcp/3456")
74 + if err != nil {
75 + t.Fatal("error setting up peer", err)
76 + }
77 +
78 + c, err := Dial("tcp", p2, p1)
79 if err != nil {
80 t.Fatal("error dialing peer", err)
81 }
net/swarm/conn.go
+44 -16
@@ -3,6 +3,8 @@ package swarm
3 import (
4 "errors"
5 "fmt"
6 + "net"
7 + "syscall"
8
9 spipe "github.com/jbenet/go-ipfs/crypto/spipe"
10 conn "github.com/jbenet/go-ipfs/net/conn"
@@ -44,6 +46,11 @@ func (s *Swarm) connListen(maddr ma.Multiaddr) error {
46 return err
47 }
48
49 + // make sure port can be reused. TOOD this doesn't work...
50 + // if err := setSocketReuse(list); err != nil {
51 + // return err
52 + // }
53 +
54 // NOTE: this may require a lock around it later. currently, only run on setup
55 s.listeners = append(s.listeners, list)
56
@@ -71,11 +78,9 @@ func (s *Swarm) connListen(maddr ma.Multiaddr) error {
78 // Handle getting ID from this peer, handshake, and adding it into the map
79 func (s *Swarm) handleIncomingConn(nconn manet.Conn) {
80
74 - addr := nconn.RemoteMultiaddr()
75 -
81 // Construct conn with nil peer for now, because we don't know its ID yet.
82 // connSetup will figure this out, and pull out / construct the peer.
78 - c, err := conn.NewConn(nil, addr, nconn)
83 + c, err := conn.NewConn(s.local, nil, nconn)
84 if err != nil {
85 s.errChan <- err
86 return
@@ -96,20 +101,20 @@ func (s *Swarm) connSetup(c *conn.Conn) error {
101 return errors.New("Tried to start nil connection.")
102 }
103
99 - if c.Peer != nil {
100 - log.Debug("Starting connection: %s", c.Peer)
104 + if c.Remote != nil {
105 + log.Debug("%s Starting connection: %s", c.Local, c.Remote)
106 } else {
102 - log.Debug("Starting connection: [unknown peer]")
107 + log.Debug("%s Starting connection: [unknown peer]", c.Local)
108 }
109
110 if err := s.connSecure(c); err != nil {
111 return fmt.Errorf("Conn securing error: %v", err)
112 }
113
109 - log.Debug("Secured connection: %s", c.Peer)
114 + log.Debug("%s secured connection: %s", c.Local, c.Remote)
115
116 // add address of connection to Peer. Maybe it should happen in connSecure.
112 - c.Peer.AddAddress(c.Addr)
117 + c.Remote.AddAddress(c.Conn.RemoteMultiaddr())
118
119 if err := s.connVersionExchange(c); err != nil {
120 return fmt.Errorf("Conn version exchange error: %v", err)
@@ -117,12 +122,12 @@ func (s *Swarm) connSetup(c *conn.Conn) error {
122
123 // add to conns
124 s.connsLock.Lock()
120 - if _, ok := s.conns[c.Peer.Key()]; ok {
125 + if _, ok := s.conns[c.Remote.Key()]; ok {
126 log.Debug("Conn already open!")
127 s.connsLock.Unlock()
128 return ErrAlreadyOpen
129 }
125 - s.conns[c.Peer.Key()] = c
130 + s.conns[c.Remote.Key()] = c
131 log.Debug("Added conn to map!")
132 s.connsLock.Unlock()
133
@@ -147,10 +152,10 @@ func (s *Swarm) connSecure(c *conn.Conn) error {
152 return err
153 }
154
150 - if c.Peer == nil {
151 - c.Peer = sp.RemotePeer()
155 + if c.Remote == nil {
156 + c.Remote = sp.RemotePeer()
157
153 - } else if c.Peer != sp.RemotePeer() {
158 + } else if c.Remote != sp.RemotePeer() {
159 panic("peers not being constructed correctly.")
160 }
161
@@ -251,20 +256,43 @@ func (s *Swarm) fanIn(c *conn.Conn) {
256
257 case data, ok := <-c.Secure.In:
258 if !ok {
254 - e := fmt.Errorf("Error retrieving from conn: %v", c.Peer)
259 + e := fmt.Errorf("Error retrieving from conn: %v", c.Remote)
260 s.errChan <- e
261 goto out
262 }
263
264 // log.Debug("[peer: %s] Received message [from = %s]", s.local, c.Peer)
265
261 - msg := msg.New(c.Peer, data)
266 + msg := msg.New(c.Remote, data)
267 s.Incoming <- msg
268 }
269 }
270
271 out:
272 s.connsLock.Lock()
268 - delete(s.conns, c.Peer.Key())
273 + delete(s.conns, c.Remote.Key())
274 s.connsLock.Unlock()
275 }
276 +
277 +func setSocketReuse(l manet.Listener) error {
278 + nl := l.NetListener()
279 +
280 + // for now only TCP. TODO change this when more networks.
281 + file, err := nl.(*net.TCPListener).File()
282 + if err != nil {
283 + return err
284 + }
285 +
286 + fd := file.Fd()
287 + err = syscall.SetsockoptInt(int(fd), syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1)
288 + if err != nil {
289 + return err
290 + }
291 +
292 + err = syscall.SetsockoptInt(int(fd), syscall.SOL_SOCKET, syscall.SO_REUSEPORT, 1)
293 + if err != nil {
294 + return err
295 + }
296 +
297 + return nil
298 +}
net/swarm/swarm.go
+4 -3
@@ -129,7 +129,7 @@ func (s *Swarm) Dial(peer *peer.Peer) (*conn.Conn, error) {
129 }
130
131 // open connection to peer
132 - c, err = conn.Dial("tcp", peer)
132 + c, err = conn.Dial("tcp", s.local, peer)
133 if err != nil {
134 return nil, err
135 }
@@ -153,7 +153,7 @@ func (s *Swarm) DialAddr(addr ma.Multiaddr) (*conn.Conn, error) {
153 npeer := new(peer.Peer)
154 npeer.AddAddress(addr)
155
156 - c, err := conn.Dial("tcp", npeer)
156 + c, err := conn.Dial("tcp", s.local, npeer)
157 if err != nil {
158 return nil, err
159 }
@@ -201,11 +201,12 @@ func (s *Swarm) GetErrChan() chan error {
201 return s.errChan
202 }
203
204 +// GetPeerList returns a copy of the set of peers swarm is connected to.
205 func (s *Swarm) GetPeerList() []*peer.Peer {
206 var out []*peer.Peer
207 s.connsLock.RLock()
208 for _, p := range s.conns {
208 - out = append(out, p.Peer)
209 + out = append(out, p.Remote)
210 }
211 s.connsLock.RUnlock()
212 return out