switching to separated conn and interface-conn packages

This commit is contained in:
Jeromy 2016-10-04 15:39:24 -07:00
parent a54a624e4a
commit 9910e6a7cf
7 changed files with 36 additions and 30 deletions

View File

@ -7,11 +7,11 @@ import (
peer "github.com/ipfs/go-libp2p-peer"
ma "github.com/jbenet/go-multiaddr"
addrutil "github.com/libp2p/go-addr-util"
conn "github.com/libp2p/go-libp2p-conn"
iconn "github.com/libp2p/go-libp2p-interface-conn"
)
type dialResult struct {
Conn conn.Conn
Conn iconn.Conn
Err error
}
@ -38,14 +38,14 @@ type dialLimiter struct {
fdLimit int
waitingOnFd []*dialJob
dialFunc func(context.Context, peer.ID, ma.Multiaddr) (conn.Conn, error)
dialFunc func(context.Context, peer.ID, ma.Multiaddr) (iconn.Conn, error)
activePerPeer map[peer.ID]int
perPeerLimit int
waitingOnPeerLimit map[peer.ID][]*dialJob
}
type dialfunc func(context.Context, peer.ID, ma.Multiaddr) (conn.Conn, error)
type dialfunc func(context.Context, peer.ID, ma.Multiaddr) (iconn.Conn, error)
func newDialLimiter(df dialfunc) *dialLimiter {
return newDialLimiterWithParams(df, concurrentFdDials, defaultPerPeerRateLimit)

View File

@ -10,7 +10,7 @@ import (
peer "github.com/ipfs/go-libp2p-peer"
ma "github.com/jbenet/go-multiaddr"
conn "github.com/libp2p/go-libp2p-conn"
iconn "github.com/libp2p/go-libp2p-interface-conn"
mafmt "github.com/whyrusleeping/mafmt"
)
@ -55,13 +55,13 @@ func tryDialAddrs(ctx context.Context, l *dialLimiter, p peer.ID, addrs []ma.Mul
}
func hangDialFunc(hang chan struct{}) dialfunc {
return func(ctx context.Context, p peer.ID, a ma.Multiaddr) (conn.Conn, error) {
return func(ctx context.Context, p peer.ID, a ma.Multiaddr) (iconn.Conn, error) {
if mafmt.UTP.Matches(a) {
return conn.Conn(nil), nil
return iconn.Conn(nil), nil
}
if tcpPortOver(a, 10) {
return conn.Conn(nil), nil
return iconn.Conn(nil), nil
}
<-hang
@ -171,9 +171,9 @@ func TestFDLimiting(t *testing.T) {
func TestTokenRedistribution(t *testing.T) {
hangchs := make(map[peer.ID]chan struct{})
df := func(ctx context.Context, p peer.ID, a ma.Multiaddr) (conn.Conn, error) {
df := func(ctx context.Context, p peer.ID, a ma.Multiaddr) (iconn.Conn, error) {
if tcpPortOver(a, 10) {
return (conn.Conn)(nil), nil
return (iconn.Conn)(nil), nil
}
<-hangchs[p]
@ -260,9 +260,9 @@ func TestTokenRedistribution(t *testing.T) {
}
func TestStressLimiter(t *testing.T) {
df := func(ctx context.Context, p peer.ID, a ma.Multiaddr) (conn.Conn, error) {
df := func(ctx context.Context, p peer.ID, a ma.Multiaddr) (iconn.Conn, error) {
if tcpPortOver(a, 1000) {
return conn.Conn(nil), nil
return iconn.Conn(nil), nil
}
time.Sleep(time.Millisecond * time.Duration(5+rand.Intn(100)))

View File

@ -3,7 +3,7 @@ package swarm
import (
ma "github.com/jbenet/go-multiaddr"
addrutil "github.com/libp2p/go-addr-util"
conn "github.com/libp2p/go-libp2p-conn"
iconn "github.com/libp2p/go-libp2p-interface-conn"
)
// ListenAddresses returns a list of addresses at which this swarm listens.
@ -11,7 +11,7 @@ func (s *Swarm) ListenAddresses() []ma.Multiaddr {
listeners := s.swarm.Listeners()
addrs := make([]ma.Multiaddr, 0, len(listeners))
for _, l := range listeners {
if l2, ok := l.NetListener().(conn.Listener); ok {
if l2, ok := l.NetListener().(iconn.Listener); ok {
addrs = append(addrs, l2.Multiaddr())
}
}

View File

@ -4,13 +4,12 @@ import (
"context"
"fmt"
inet "github.com/libp2p/go-libp2p-net"
ic "github.com/ipfs/go-libp2p-crypto"
peer "github.com/ipfs/go-libp2p-peer"
ma "github.com/jbenet/go-multiaddr"
ps "github.com/jbenet/go-peerstream"
conn "github.com/libp2p/go-libp2p-conn"
iconn "github.com/libp2p/go-libp2p-interface-conn"
inet "github.com/libp2p/go-libp2p-net"
)
// Conn is a simple wrapper around a ps.Conn that also exposes
@ -33,12 +32,12 @@ func (c *Conn) StreamConn() *ps.Conn {
return (*ps.Conn)(c)
}
func (c *Conn) RawConn() conn.Conn {
func (c *Conn) RawConn() iconn.Conn {
// righly panic if these things aren't true. it is an expected
// invariant that these Conns are all of the typewe expect:
// ps.Conn wrapping a conn.Conn
// if we get something else it is programmer error.
return (*ps.Conn)(c).NetConn().(conn.Conn)
return (*ps.Conn)(c).NetConn().(iconn.Conn)
}
func (c *Conn) String() string {
@ -94,7 +93,7 @@ func (c *Conn) Close() error {
func wrapConn(psc *ps.Conn) (*Conn, error) {
// grab the underlying connection.
if _, ok := psc.NetConn().(conn.Conn); !ok {
if _, ok := psc.NetConn().(iconn.Conn); !ok {
// this should never happen. if we see it ocurring it means that we added
// a Listener to the ps.Swarm that is NOT one of our net/conn.Listener.
return nil, fmt.Errorf("swarm connHandler: invalid conn (not a conn.Conn): %s", psc)

View File

@ -11,7 +11,7 @@ import (
peer "github.com/ipfs/go-libp2p-peer"
ma "github.com/jbenet/go-multiaddr"
addrutil "github.com/libp2p/go-addr-util"
conn "github.com/libp2p/go-libp2p-conn"
iconn "github.com/libp2p/go-libp2p-interface-conn"
)
// Diagram of dial sync:
@ -284,7 +284,7 @@ func (s *Swarm) dial(ctx context.Context, p peer.ID) (*Conn, error) {
return swarmC, nil
}
func (s *Swarm) dialAddrs(ctx context.Context, p peer.ID, remoteAddrs <-chan ma.Multiaddr) (conn.Conn, error) {
func (s *Swarm) dialAddrs(ctx context.Context, p peer.ID, remoteAddrs <-chan ma.Multiaddr) (iconn.Conn, error) {
log.Debugf("%s swarm dialing %s %s", s.local, p, remoteAddrs)
ctx, cancel := context.WithCancel(ctx)
@ -344,7 +344,7 @@ func (s *Swarm) limitedDial(ctx context.Context, p peer.ID, a ma.Multiaddr, resp
})
}
func (s *Swarm) dialAddr(ctx context.Context, p peer.ID, addr ma.Multiaddr) (conn.Conn, error) {
func (s *Swarm) dialAddr(ctx context.Context, p peer.ID, addr ma.Multiaddr) (iconn.Conn, error) {
log.Debugf("%s swarm dialing %s %s", s.local, p, addr)
connC, err := s.dialer.Dial(ctx, addr, p)
@ -376,7 +376,7 @@ var ConnSetupTimeout = time.Minute * 5
// dialConnSetup is the setup logic for a connection from the dial side. it
// needs to add the Conn to the StreamSwarm, then run newConnSetup
func dialConnSetup(ctx context.Context, s *Swarm, connC conn.Conn) (*Conn, error) {
func dialConnSetup(ctx context.Context, s *Swarm, connC iconn.Conn) (*Conn, error) {
deadline, ok := ctx.Deadline()
if !ok {

View File

@ -8,6 +8,7 @@ import (
ma "github.com/jbenet/go-multiaddr"
ps "github.com/jbenet/go-peerstream"
conn "github.com/libp2p/go-libp2p-conn"
iconn "github.com/libp2p/go-libp2p-interface-conn"
mconn "github.com/libp2p/go-libp2p-metrics/conn"
inet "github.com/libp2p/go-libp2p-net"
transport "github.com/libp2p/go-libp2p-transport"
@ -98,7 +99,7 @@ func (s *Swarm) addListener(tptlist transport.Listener) error {
return s.addConnListener(list)
}
func (s *Swarm) addConnListener(list conn.Listener) error {
func (s *Swarm) addConnListener(list iconn.Listener) error {
// AddListener to the peerstream Listener. this will begin accepting connections
// and streams!
sl, err := s.swarm.AddListener(list)

View File

@ -215,21 +215,27 @@
},
{
"author": "whyrusleeping",
"hash": "QmTaW4q1AbqMkpfDLUYzW18nW62GsrnFvtVcvR1pnaURm6",
"hash": "Qmd4eGc7VMmVZu225mdCzsagdWJBaAVHmmvNFYSvxSUeX5",
"name": "go-libp2p-conn",
"version": "1.0.0"
"version": "1.1.0"
},
{
"author": "whyrusleeping",
"hash": "QmSctCnUNkE7c7C2LGRGYrdmU9SDH3MtbAeueN7Jq1NN2q",
"hash": "Qmcr6YS5QCPwejhBxZk5iFEvVEHSVzewbLVMyc64zefaUf",
"name": "go-libp2p-net",
"version": "1.0.0"
},
{
"author": "whyrusleeping",
"hash": "QmS38TTyj47fP5NVNznxpw4KpFYSU9zHkCEfmuBSkHzb6d",
"hash": "QmbidKjoAxLVE2skY74Rfw1QafReyf4w2FZqA8Sd7MxoVd",
"name": "go-libp2p-metrics",
"version": "1.0.0"
"version": "1.1.0"
},
{
"author": "whyrusleeping",
"hash": "QmRxzoGdQaN6HYyqWnT82NnLLHBAZbTUvxPEfTBTjU7KCn",
"name": "go-libp2p-interface-conn",
"version": "0.1.0"
}
],
"gxVersion": "0.4.0",