@cryptotaxi247 / kubo / commits / 8e6609e78

add timeout opt to transport dialer creation

License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>

Jeromy committed Nov 3, 2015 at 09:48 UTC 8e6609e78a0dc9e5307ae9d0b7502040301507ab
5 files changed +28 -8
p2p/net/conn/dial.go
+1 -2
@@ -4,7 +4,6 @@ import (
4 "fmt"
5 "math/rand"
6 "strings"
7 - "time"
7
8 ma "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
9 manet "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net"
@@ -19,7 +18,7 @@ import (
18
19 type WrapFunc func(transport.Conn) transport.Conn
20
22 -func NewDialer(p peer.ID, pk ci.PrivKey, tout time.Duration, wrap WrapFunc) *Dialer {
21 +func NewDialer(p peer.ID, pk ci.PrivKey, wrap WrapFunc) *Dialer {
22 return &Dialer{
23 LocalPeer: p,
24 PrivateKey: pk,
p2p/net/swarm/swarm.go
+1 -1
@@ -103,7 +103,7 @@ func NewSwarm(ctx context.Context, listenAddrs []ma.Multiaddr,
103 bwc: bwc,
104 fdRateLimit: make(chan struct{}, concurrentFdDials),
105 Filters: filter.NewFilters(),
106 - dialer: conn.NewDialer(local, peers.PrivKey(local), DialTimeout, wrap),
106 + dialer: conn.NewDialer(local, peers.PrivKey(local), wrap),
107 }
108
109 // configure Swarm
p2p/net/swarm/swarm_listen.go
+1 -1
@@ -22,7 +22,7 @@ func (s *Swarm) setupAddresses(addrs []ma.Multiaddr) error {
22 return fmt.Errorf("no transport for address: %s", a)
23 }
24
25 - d, err := tpt.Dialer(a)
25 + d, err := tpt.Dialer(a, transport.TimeoutOpt(DialTimeout))
26 if err != nil {
27 return err
28 }
p2p/net/transport/tcp.go
+21 -3
@@ -1,8 +1,10 @@
1 package transport
2
3 import (
4 + "fmt"
5 "net"
6 "sync"
7 + "time"
8
9 ma "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
10 manet "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr-net"
@@ -26,7 +28,7 @@ func NewTCPTransport() *TcpTransport {
28 }
29 }
30
29 -func (t *TcpTransport) Dialer(laddr ma.Multiaddr) (Dialer, error) {
31 +func (t *TcpTransport) Dialer(laddr ma.Multiaddr, opts ...DialOpt) (Dialer, error) {
32 t.dlock.Lock()
33 defer t.dlock.Unlock()
34 s := laddr.String()
@@ -34,8 +36,17 @@ func (t *TcpTransport) Dialer(laddr ma.Multiaddr) (Dialer, error) {
36 if found {
37 return d, nil
38 }
37 -
39 var base manet.Dialer
40 +
41 + for _, o := range opts {
42 + switch o := o.(type) {
43 + case TimeoutOpt:
44 + base.Timeout = o.(time.Duration)
45 + default:
46 + return nil, fmt.Errorf("unrecognized option: %#v", o)
47 + }
48 + }
49 +
50 tcpd, err := t.newTcpDialer(base, laddr)
51 if err != nil {
52 return nil, err
@@ -119,10 +130,17 @@ func (t *TcpTransport) newTcpDialer(base manet.Dialer, laddr ma.Multiaddr) (*tcp
130 }, nil
131 }
132
133 + rd := reuseport.Dialer{
134 + D: net.Dialer{
135 + LocalAddr: la,
136 + Timeout: base.Timeout,
137 + },
138 + }
139 +
140 return &tcpDialer{
141 doReuse: true,
142 laddr: laddr,
125 - rd: reuseport.Dialer{D: net.Dialer{LocalAddr: la}},
143 + rd: rd,
144 madialer: base,
145 transport: t,
146 }, nil
p2p/net/transport/transport.go
+4 -1
@@ -17,7 +17,7 @@ type Conn interface {
17 }
18
19 type Transport interface {
20 - Dialer(laddr ma.Multiaddr) (Dialer, error)
20 + Dialer(laddr ma.Multiaddr, opts ...DialOpt) (Dialer, error)
21 Listener(laddr ma.Multiaddr) (Listener, error)
22 Matches(ma.Multiaddr) bool
23 }
@@ -43,6 +43,9 @@ func (cw *connWrap) Transport() Transport {
43 return cw.transport
44 }
45
46 +type DialOpt interface{}
47 +type TimeoutOpt interface{}
48 +
49 func IsTcpMultiaddr(a ma.Multiaddr) bool {
50 p := a.Protocols()
51 return len(p) == 2 && (p[0].Name == "ip4" || p[0].Name == "ip6") && p[1].Name == "tcp"