@cryptotaxi247 / kubo / commits / 724e6cb37

reuseport tcp is now a dial creation option

And other assorted PR feedback License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Nov 8, 2015 at 11:04 UTC 724e6cb37d8a32835d50b98543b6df50ce44d6a0
5 files changed +30 -23
p2p/net/conn/dial_test.go
+1 -1
@@ -54,7 +54,7 @@ func setupSingleConn(t *testing.T, ctx context.Context) (a, b Conn, p1, p2 tu.Pe
54 }
55
56 func Listen(ctx context.Context, addr ma.Multiaddr, local peer.ID, sk ic.PrivKey) (Listener, error) {
57 - list, err := transport.NewTCPTransport().Listener(addr)
57 + list, err := transport.NewTCPTransport().Listen(addr)
58 if err != nil {
59 return nil, err
60 }
p2p/net/swarm/swarm.go
+2 -2
@@ -114,7 +114,7 @@ func NewSwarm(ctx context.Context, listenAddrs []ma.Multiaddr,
114 prom.MustRegisterOrGet(peersTotal)
115 s.Notify((*metricsNotifiee)(s))
116
117 - return s, s.setupAddresses(listenAddrs)
117 + return s, s.setupInterfaces(listenAddrs)
118 }
119
120 func (s *Swarm) teardown() error {
@@ -147,7 +147,7 @@ func (s *Swarm) Listen(addrs ...ma.Multiaddr) error {
147 return err
148 }
149
150 - return s.setupAddresses(addrs)
150 + return s.setupInterfaces(addrs)
151 }
152
153 // Process returns the Process of the swarm
p2p/net/swarm/swarm_listen.go
+3 -3
@@ -15,21 +15,21 @@ import (
15 )
16
17 // Open listeners and reuse-dialers for the given addresses
18 -func (s *Swarm) setupAddresses(addrs []ma.Multiaddr) error {
18 +func (s *Swarm) setupInterfaces(addrs []ma.Multiaddr) error {
19 for _, a := range addrs {
20 tpt := s.transportForAddr(a)
21 if tpt == nil {
22 return fmt.Errorf("no transport for address: %s", a)
23 }
24
25 - d, err := tpt.Dialer(a, transport.TimeoutOpt(DialTimeout))
25 + d, err := tpt.Dialer(a, transport.TimeoutOpt(DialTimeout), transport.ReusePorts)
26 if err != nil {
27 return err
28 }
29
30 s.dialer.AddDialer(d)
31
32 - list, err := tpt.Listener(a)
32 + list, err := tpt.Listen(a)
33 if err != nil {
34 return err
35 }
p2p/net/transport/tcp.go
+18 -15
@@ -38,16 +38,19 @@ func (t *TcpTransport) Dialer(laddr ma.Multiaddr, opts ...DialOpt) (Dialer, erro
38 }
39 var base manet.Dialer
40
41 + var doReuse bool
42 for _, o := range opts {
43 switch o := o.(type) {
44 case TimeoutOpt:
44 - base.Timeout = o.(time.Duration)
45 + base.Timeout = time.Duration(o)
46 + case ReuseportOpt:
47 + doReuse = bool(o)
48 default:
49 return nil, fmt.Errorf("unrecognized option: %#v", o)
50 }
51 }
52
50 - tcpd, err := t.newTcpDialer(base, laddr)
53 + tcpd, err := t.newTcpDialer(base, laddr, doReuse)
54 if err != nil {
55 return nil, err
56 }
@@ -56,7 +59,7 @@ func (t *TcpTransport) Dialer(laddr ma.Multiaddr, opts ...DialOpt) (Dialer, erro
59 return tcpd, nil
60 }
61
59 -func (t *TcpTransport) Listener(laddr ma.Multiaddr) (Listener, error) {
62 +func (t *TcpTransport) Listen(laddr ma.Multiaddr) (Listener, error) {
63 t.llock.Lock()
64 defer t.llock.Unlock()
65 s := laddr.String()
@@ -114,33 +117,33 @@ type tcpDialer struct {
117 transport Transport
118 }
119
117 -func (t *TcpTransport) newTcpDialer(base manet.Dialer, laddr ma.Multiaddr) (*tcpDialer, error) {
120 +func (t *TcpTransport) newTcpDialer(base manet.Dialer, laddr ma.Multiaddr, doReuse bool) (*tcpDialer, error) {
121 // get the local net.Addr manually
122 la, err := manet.ToNetAddr(laddr)
123 if err != nil {
124 return nil, err // something wrong with laddr.
125 }
126
124 - if !ReuseportIsAvailable() {
127 + if doReuse && ReuseportIsAvailable() {
128 + rd := reuseport.Dialer{
129 + D: net.Dialer{
130 + LocalAddr: la,
131 + Timeout: base.Timeout,
132 + },
133 + }
134 +
135 return &tcpDialer{
126 - doReuse: false,
136 + doReuse: true,
137 laddr: laddr,
138 + rd: rd,
139 madialer: base,
140 transport: t,
141 }, nil
142 }
143
133 - rd := reuseport.Dialer{
134 - D: net.Dialer{
135 - LocalAddr: la,
136 - Timeout: base.Timeout,
137 - },
138 - }
139 -
144 return &tcpDialer{
141 - doReuse: true,
145 + doReuse: false,
146 laddr: laddr,
143 - rd: rd,
147 madialer: base,
148 transport: t,
149 }, nil
p2p/net/transport/transport.go
+6 -2
@@ -2,6 +2,7 @@ package transport
2
3 import (
4 "net"
5 + "time"
6
7 ma "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
8 manet "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net"
@@ -18,7 +19,7 @@ type Conn interface {
19
20 type Transport interface {
21 Dialer(laddr ma.Multiaddr, opts ...DialOpt) (Dialer, error)
21 - Listener(laddr ma.Multiaddr) (Listener, error)
22 + Listen(laddr ma.Multiaddr) (Listener, error)
23 Matches(ma.Multiaddr) bool
24 }
25
@@ -44,7 +45,10 @@ func (cw *connWrap) Transport() Transport {
45 }
46
47 type DialOpt interface{}
47 -type TimeoutOpt interface{}
48 +type TimeoutOpt time.Duration
49 +type ReuseportOpt bool
50 +
51 +var ReusePorts ReuseportOpt = true
52
53 func IsTcpMultiaddr(a ma.Multiaddr) bool {
54 p := a.Protocols()