mirror of
https://github.com/netbirdio/ice.git
synced 2026-05-22 17:10:58 -07:00
Add vnet support with the new pion/turn client
Relates to pion/webrtc#712
This commit is contained in:
committed by
Yutaka Takeda
parent
d2d97a6293
commit
b3e05469f0
@@ -16,6 +16,7 @@ import (
|
||||
"github.com/pion/mdns"
|
||||
"github.com/pion/stun"
|
||||
"github.com/pion/transport/packetio"
|
||||
"github.com/pion/transport/vnet"
|
||||
"golang.org/x/net/ipv4"
|
||||
)
|
||||
|
||||
@@ -141,7 +142,10 @@ type Agent struct {
|
||||
done chan struct{}
|
||||
err atomicError
|
||||
|
||||
log logging.LeveledLogger
|
||||
loggerFactory logging.LoggerFactory
|
||||
log logging.LeveledLogger
|
||||
|
||||
net *vnet.Net
|
||||
}
|
||||
|
||||
func (a *Agent) ok() error {
|
||||
@@ -219,6 +223,10 @@ type AgentConfig struct {
|
||||
PrflxAcceptanceMinWait *time.Duration
|
||||
// HostAcceptanceMinWait specify a minimum wait time before selecting relay candidates
|
||||
RelayAcceptanceMinWait *time.Duration
|
||||
|
||||
// Net is the our abstracted network interface for internal development purpose only
|
||||
// (see github.com/pion/transport/vnet)
|
||||
Net *vnet.Net
|
||||
}
|
||||
|
||||
// NewAgent creates a new Agent
|
||||
@@ -278,16 +286,18 @@ func NewAgent(config *AgentConfig) (*Agent, error) {
|
||||
urls: config.Urls,
|
||||
networkTypes: config.NetworkTypes,
|
||||
|
||||
localUfrag: randSeq(16),
|
||||
localPwd: randSeq(32),
|
||||
taskChan: make(chan task),
|
||||
onConnected: make(chan struct{}),
|
||||
buffer: packetio.NewBuffer(),
|
||||
done: make(chan struct{}),
|
||||
portmin: config.PortMin,
|
||||
portmax: config.PortMax,
|
||||
trickle: config.Trickle,
|
||||
log: loggerFactory.NewLogger("ice"),
|
||||
localUfrag: randSeq(16),
|
||||
localPwd: randSeq(32),
|
||||
taskChan: make(chan task),
|
||||
onConnected: make(chan struct{}),
|
||||
buffer: packetio.NewBuffer(),
|
||||
done: make(chan struct{}),
|
||||
portmin: config.PortMin,
|
||||
portmax: config.PortMax,
|
||||
trickle: config.Trickle,
|
||||
loggerFactory: loggerFactory,
|
||||
log: loggerFactory.NewLogger("ice"),
|
||||
net: config.Net,
|
||||
|
||||
mDNSMode: mDNSMode,
|
||||
mDNSName: mDNSName,
|
||||
@@ -297,6 +307,15 @@ func NewAgent(config *AgentConfig) (*Agent, error) {
|
||||
}
|
||||
a.haveStarted.Store(false)
|
||||
|
||||
if a.net == nil {
|
||||
a.net = vnet.NewNet(nil)
|
||||
} else {
|
||||
a.log.Warn("vnet is enabled")
|
||||
if a.mDNSMode != MulticastDNSModeDisabled {
|
||||
a.log.Warn("vnet does not support mDNS yet")
|
||||
}
|
||||
}
|
||||
|
||||
if config.MaxBindingRequests == nil {
|
||||
a.maxBindingRequests = defaultMaxBindingRequests
|
||||
} else {
|
||||
@@ -675,14 +694,6 @@ func (a *Agent) addRemoteCandidate(c Candidate) {
|
||||
set = append(set, c)
|
||||
a.remoteCandidates[c.NetworkType()] = set
|
||||
|
||||
for _, l := range a.localCandidates[NetworkTypeUDP4] {
|
||||
if localRelay, ok := l.(*CandidateRelay); ok {
|
||||
if err := localRelay.addPermission(c); err != nil {
|
||||
a.log.Errorf("Failed to create TURN permission %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if localCandidates, ok := a.localCandidates[c.NetworkType()]; ok {
|
||||
for _, localCandidate := range localCandidates {
|
||||
a.addPair(localCandidate, c)
|
||||
|
||||
+8
-64
@@ -1,20 +1,14 @@
|
||||
package ice
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"io"
|
||||
"net"
|
||||
|
||||
"github.com/pion/turnc"
|
||||
)
|
||||
|
||||
// CandidateRelay ...
|
||||
type CandidateRelay struct {
|
||||
candidateBase
|
||||
|
||||
allocation *turnc.Allocation
|
||||
client io.Closer
|
||||
permissions map[string]*turnc.Permission
|
||||
onClose func() error
|
||||
}
|
||||
|
||||
// CandidateRelayConfig is the config required to create a new CandidateRelay
|
||||
@@ -26,6 +20,7 @@ type CandidateRelayConfig struct {
|
||||
Component uint16
|
||||
RelAddr string
|
||||
RelPort int
|
||||
OnClose func() error
|
||||
}
|
||||
|
||||
// NewCandidateRelay creates a new relay candidate
|
||||
@@ -64,66 +59,15 @@ func NewCandidateRelay(config *CandidateRelayConfig) (*CandidateRelay, error) {
|
||||
Port: config.RelPort,
|
||||
},
|
||||
},
|
||||
permissions: map[string]*turnc.Permission{},
|
||||
onClose: config.OnClose,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (c *CandidateRelay) setAllocation(client io.Closer, a *turnc.Allocation) {
|
||||
c.allocation = a
|
||||
c.client = client
|
||||
}
|
||||
|
||||
func (c *CandidateRelay) start(a *Agent, conn net.PacketConn) {
|
||||
c.currAgent = a
|
||||
}
|
||||
|
||||
func (c *CandidateRelay) close() error {
|
||||
c.lock.Lock()
|
||||
defer c.lock.Unlock()
|
||||
for _, p := range c.permissions {
|
||||
if err := p.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
err := c.candidateBase.close()
|
||||
if c.onClose != nil {
|
||||
c.onClose()
|
||||
c.onClose = nil
|
||||
}
|
||||
if c.client == nil {
|
||||
return nil
|
||||
}
|
||||
return c.client.Close()
|
||||
}
|
||||
|
||||
func (c *CandidateRelay) addPermission(dst Candidate) error {
|
||||
permission, err := c.allocation.Create(dst.addr())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
c.lock.Lock()
|
||||
c.permissions[dst.String()] = permission
|
||||
if err = c.permissions[dst.String()].Bind(); err != nil {
|
||||
c.agent().log.Warnf("Failed to Create ChannelBind for %v: %v", dst.String, err)
|
||||
}
|
||||
c.lock.Unlock()
|
||||
|
||||
go func(remoteAddr net.Addr) {
|
||||
log := c.agent().log
|
||||
buffer := make([]byte, receiveMTU)
|
||||
for {
|
||||
n, err := permission.Read(buffer)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
handleInboundCandidateMsg(c, buffer[:n], remoteAddr, log)
|
||||
}
|
||||
}(dst.addr())
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *CandidateRelay) writeTo(raw []byte, dst Candidate) (int, error) {
|
||||
permission, ok := c.permissions[dst.String()]
|
||||
if !ok {
|
||||
return 0, errors.New("no permission created for remote candidate")
|
||||
}
|
||||
|
||||
return permission.Write(raw)
|
||||
return err
|
||||
}
|
||||
|
||||
+19
-14
@@ -5,14 +5,12 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pion/logging"
|
||||
"github.com/pion/transport/test"
|
||||
"github.com/pion/turn"
|
||||
)
|
||||
|
||||
type mockTURNServer struct {
|
||||
}
|
||||
|
||||
func (m *mockTURNServer) AuthenticateRequest(username string, srcAddr net.Addr) (password string, ok bool) {
|
||||
func optimisticAuthHandler(username string, srcAddr net.Addr) (password string, ok bool) {
|
||||
return "password", true
|
||||
}
|
||||
|
||||
@@ -24,16 +22,19 @@ func TestRelayOnlyConnection(t *testing.T) {
|
||||
report := test.CheckRoutines(t)
|
||||
defer report()
|
||||
|
||||
serverPort := randomPort(t)
|
||||
server := turn.Create(turn.StartArguments{
|
||||
Server: &mockTURNServer{},
|
||||
Realm: "localhost",
|
||||
})
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
|
||||
serverChan := make(chan error, 1)
|
||||
go func() {
|
||||
serverChan <- server.Listen("0.0.0.0", serverPort)
|
||||
}()
|
||||
serverPort := randomPort(t)
|
||||
server := turn.NewServer(&turn.ServerConfig{
|
||||
Realm: "pion.ly",
|
||||
AuthHandler: optimisticAuthHandler,
|
||||
ListeningPort: serverPort,
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
err := server.AddListeningIPAddr("127.0.0.1")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
cfg := &AgentConfig{
|
||||
NetworkTypes: supportedNetworkTypes,
|
||||
@@ -49,6 +50,11 @@ func TestRelayOnlyConnection(t *testing.T) {
|
||||
CandidateTypes: []CandidateType{CandidateTypeRelay},
|
||||
}
|
||||
|
||||
err = server.Start()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
aAgent, err := NewAgent(cfg)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -82,5 +88,4 @@ func TestRelayOnlyConnection(t *testing.T) {
|
||||
if err = server.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
<-serverChan
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pion/logging"
|
||||
"github.com/pion/transport/test"
|
||||
"github.com/pion/turn"
|
||||
)
|
||||
@@ -16,16 +17,25 @@ func TestServerReflexiveOnlyConnection(t *testing.T) {
|
||||
report := test.CheckRoutines(t)
|
||||
defer report()
|
||||
|
||||
serverPort := randomPort(t)
|
||||
server := turn.Create(turn.StartArguments{
|
||||
Server: &mockTURNServer{},
|
||||
Realm: "localhost",
|
||||
})
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
//log := loggerFactory.NewLogger("test")
|
||||
|
||||
serverChan := make(chan error, 1)
|
||||
go func() {
|
||||
serverChan <- server.Listen("0.0.0.0", serverPort)
|
||||
}()
|
||||
serverPort := randomPort(t)
|
||||
server := turn.NewServer(&turn.ServerConfig{
|
||||
Realm: "pion.ly",
|
||||
AuthHandler: optimisticAuthHandler,
|
||||
ListeningPort: serverPort,
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
err := server.AddListeningIPAddr("127.0.0.1")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
err = server.Start()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
cfg := &AgentConfig{
|
||||
NetworkTypes: []NetworkType{NetworkTypeUDP4},
|
||||
@@ -72,5 +82,4 @@ func TestServerReflexiveOnlyConnection(t *testing.T) {
|
||||
if err = server.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
<-serverChan
|
||||
}
|
||||
|
||||
@@ -0,0 +1,340 @@
|
||||
package ice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/pion/logging"
|
||||
"github.com/pion/transport/vnet"
|
||||
"github.com/pion/turn"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
type virtualNet struct {
|
||||
wan *vnet.Router
|
||||
net0 *vnet.Net
|
||||
net1 *vnet.Net
|
||||
server *turn.Server
|
||||
}
|
||||
|
||||
func (v *virtualNet) close() {
|
||||
v.server.Close() // nolint:errcheck
|
||||
v.wan.Stop()
|
||||
}
|
||||
|
||||
func buildVNet(natType *vnet.NATType) (*virtualNet, error) {
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
|
||||
// WAN
|
||||
wan, err := vnet.NewRouter(&vnet.RouterConfig{
|
||||
CIDR: "0.0.0.0/0",
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
wanNet := vnet.NewNet(&vnet.NetConfig{
|
||||
StaticIP: "1.2.3.4", // will be assigned to eth0
|
||||
})
|
||||
|
||||
err = wan.AddNet(wanNet)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// LAN 0
|
||||
lan0, err := vnet.NewRouter(&vnet.RouterConfig{
|
||||
StaticIP: "27.1.1.1", // this router's external IP on eth0
|
||||
CIDR: "192.168.0.0/24",
|
||||
NATType: natType,
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
net0 := vnet.NewNet(&vnet.NetConfig{})
|
||||
err = lan0.AddNet(net0)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = wan.AddRouter(lan0)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// LAN 1
|
||||
lan1, err := vnet.NewRouter(&vnet.RouterConfig{
|
||||
StaticIP: "28.1.1.1", // this router's external IP on eth0
|
||||
CIDR: "10.2.0.0/24",
|
||||
NATType: natType,
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
net1 := vnet.NewNet(&vnet.NetConfig{})
|
||||
err = lan1.AddNet(net1)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = wan.AddRouter(lan1)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Start routers
|
||||
err = wan.Start()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Run TURN(STUN) server
|
||||
credMap := map[string]string{}
|
||||
credMap["user"] = "pass"
|
||||
server := turn.NewServer(&turn.ServerConfig{
|
||||
AuthHandler: func(username string, srcAddr net.Addr) (password string, ok bool) {
|
||||
if pw, ok := credMap[username]; ok {
|
||||
return pw, true
|
||||
}
|
||||
return "", false
|
||||
},
|
||||
Realm: "pion.ly",
|
||||
Net: wanNet,
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
err = server.AddListeningIPAddr("1.2.3.4")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = server.Start()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &virtualNet{
|
||||
wan: wan,
|
||||
net0: net0,
|
||||
net1: net1,
|
||||
server: server,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func connectWithVNet(aAgent, bAgent *Agent) (*Conn, *Conn) {
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
log := loggerFactory.NewLogger("test")
|
||||
|
||||
// Manual signaling
|
||||
aUfrag, aPwd := aAgent.GetLocalUserCredentials()
|
||||
bUfrag, bPwd := bAgent.GetLocalUserCredentials()
|
||||
|
||||
candidates, err := aAgent.GetLocalCandidates()
|
||||
check(err)
|
||||
for _, c := range candidates {
|
||||
log.Debugf("agent a candidate: %v", c.String())
|
||||
check(bAgent.AddRemoteCandidate(copyCandidate(c)))
|
||||
}
|
||||
|
||||
candidates, err = bAgent.GetLocalCandidates()
|
||||
check(err)
|
||||
for _, c := range candidates {
|
||||
log.Debugf("agent b candidate: %v", c.String())
|
||||
check(aAgent.AddRemoteCandidate(copyCandidate(c)))
|
||||
}
|
||||
|
||||
accepted := make(chan struct{})
|
||||
var aConn *Conn
|
||||
|
||||
go func() {
|
||||
var acceptErr error
|
||||
aConn, acceptErr = aAgent.Accept(context.TODO(), bUfrag, bPwd)
|
||||
check(acceptErr)
|
||||
close(accepted)
|
||||
}()
|
||||
|
||||
bConn, err := bAgent.Dial(context.TODO(), aUfrag, aPwd)
|
||||
check(err)
|
||||
|
||||
// Ensure accepted
|
||||
<-accepted
|
||||
return aConn, bConn
|
||||
}
|
||||
|
||||
func pipeWithVNet(v *virtualNet, urls0, urls1 []*URL) (*Conn, *Conn) {
|
||||
aNotifier, aConnected := onConnected()
|
||||
bNotifier, bConnected := onConnected()
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
|
||||
cfg0 := &AgentConfig{
|
||||
Urls: urls0,
|
||||
Trickle: true,
|
||||
NetworkTypes: supportedNetworkTypes,
|
||||
MulticastDNSMode: MulticastDNSModeDisabled,
|
||||
Net: v.net0,
|
||||
}
|
||||
|
||||
aAgent, err := NewAgent(cfg0)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = aAgent.OnConnectionStateChange(aNotifier)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = aAgent.OnCandidate(func(candidate Candidate) {
|
||||
if candidate == nil {
|
||||
wg.Done()
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = aAgent.GatherCandidates()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
cfg1 := &AgentConfig{
|
||||
Urls: urls1,
|
||||
Trickle: true,
|
||||
NetworkTypes: supportedNetworkTypes,
|
||||
MulticastDNSMode: MulticastDNSModeDisabled,
|
||||
Net: v.net1,
|
||||
}
|
||||
|
||||
bAgent, err := NewAgent(cfg1)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = bAgent.OnConnectionStateChange(bNotifier)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = bAgent.OnCandidate(func(candidate Candidate) {
|
||||
if candidate == nil {
|
||||
wg.Done()
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
err = bAgent.GatherCandidates()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
aConn, bConn := connectWithVNet(aAgent, bAgent)
|
||||
|
||||
// Ensure pair selected
|
||||
// Note: this assumes ConnectionStateConnected is thrown after selecting the final pair
|
||||
<-aConnected
|
||||
<-bConnected
|
||||
|
||||
return aConn, bConn
|
||||
}
|
||||
|
||||
func closePipe(t *testing.T, ca *Conn, cb *Conn) bool {
|
||||
err := ca.Close()
|
||||
if !assert.NoError(t, err, "should succeed") {
|
||||
return false
|
||||
}
|
||||
err = cb.Close()
|
||||
if !assert.NoError(t, err, "should succeed") {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func TestConnectivityVNet(t *testing.T) {
|
||||
stunServerURL := &URL{
|
||||
Scheme: SchemeTypeSTUN,
|
||||
Host: "1.2.3.4",
|
||||
Port: 3478,
|
||||
Proto: ProtoTypeUDP,
|
||||
}
|
||||
|
||||
turnServerURL := &URL{
|
||||
Scheme: SchemeTypeTURN,
|
||||
Host: "1.2.3.4",
|
||||
Port: 3478,
|
||||
Username: "user",
|
||||
Password: "pass",
|
||||
Proto: ProtoTypeUDP,
|
||||
}
|
||||
|
||||
t.Run("Full-cone NATs", func(t *testing.T) {
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
log := loggerFactory.NewLogger("test")
|
||||
|
||||
// buildVNet with Full-cone NATs
|
||||
v, err := buildVNet(&vnet.NATType{
|
||||
MappingBehavior: vnet.EndpointIndependent,
|
||||
FilteringBehavior: vnet.EndpointIndependent,
|
||||
})
|
||||
|
||||
if !assert.NoError(t, err, "should succeed") {
|
||||
return
|
||||
}
|
||||
defer v.close()
|
||||
|
||||
log.Debug("Connecting...")
|
||||
urls0 := []*URL{
|
||||
stunServerURL,
|
||||
}
|
||||
|
||||
urls1 := []*URL{
|
||||
stunServerURL,
|
||||
}
|
||||
ca, cb := pipeWithVNet(v, urls0, urls1)
|
||||
|
||||
time.Sleep(1 * time.Second)
|
||||
|
||||
log.Debug("Closing...")
|
||||
if !closePipe(t, ca, cb) {
|
||||
return
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Symmetric NATs", func(t *testing.T) {
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
log := loggerFactory.NewLogger("test")
|
||||
|
||||
// buildVNet with Symmetric NATs
|
||||
v, err := buildVNet(&vnet.NATType{
|
||||
MappingBehavior: vnet.EndpointAddrPortDependent,
|
||||
FilteringBehavior: vnet.EndpointAddrPortDependent,
|
||||
})
|
||||
|
||||
if !assert.NoError(t, err, "should succeed") {
|
||||
return
|
||||
}
|
||||
defer v.close()
|
||||
|
||||
log.Debug("Connecting...")
|
||||
urls0 := []*URL{
|
||||
stunServerURL,
|
||||
turnServerURL,
|
||||
}
|
||||
|
||||
urls1 := []*URL{
|
||||
stunServerURL,
|
||||
}
|
||||
ca, cb := pipeWithVNet(v, urls0, urls1)
|
||||
|
||||
log.Debug("Closing...")
|
||||
if !closePipe(t, ca, cb) {
|
||||
return
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -7,15 +7,16 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/pion/stun"
|
||||
"github.com/pion/turnc"
|
||||
"github.com/pion/transport/vnet"
|
||||
"github.com/pion/turn"
|
||||
)
|
||||
|
||||
const (
|
||||
stunGatherTimeout = time.Second * 5
|
||||
)
|
||||
|
||||
func localInterfaces(networkTypes []NetworkType) (ips []net.IP) {
|
||||
ifaces, err := net.Interfaces()
|
||||
func (a *Agent) localInterfaces(networkTypes []NetworkType) (ips []net.IP) {
|
||||
ifaces, err := a.net.Interfaces()
|
||||
if err != nil {
|
||||
return ips
|
||||
}
|
||||
@@ -73,9 +74,9 @@ func localInterfaces(networkTypes []NetworkType) (ips []net.IP) {
|
||||
return ips
|
||||
}
|
||||
|
||||
func listenUDP(portMax, portMin int, network string, laddr *net.UDPAddr) (*net.UDPConn, error) {
|
||||
func (a *Agent) listenUDP(portMax, portMin int, network string, laddr *net.UDPAddr) (vnet.UDPPacketConn, error) {
|
||||
if (laddr.Port != 0) || ((portMin == 0) && (portMax == 0)) {
|
||||
return net.ListenUDP(network, laddr)
|
||||
return a.net.ListenUDP(network, laddr)
|
||||
}
|
||||
var i, j int
|
||||
i = portMin
|
||||
@@ -87,10 +88,12 @@ func listenUDP(portMax, portMin int, network string, laddr *net.UDPAddr) (*net.U
|
||||
j = 0xFFFF
|
||||
}
|
||||
for i <= j {
|
||||
c, e := net.ListenUDP(network, &net.UDPAddr{IP: laddr.IP, Port: i})
|
||||
laddr = &net.UDPAddr{IP: laddr.IP, Port: i}
|
||||
c, e := a.net.ListenUDP(network, laddr)
|
||||
if e == nil {
|
||||
return c, e
|
||||
}
|
||||
a.log.Debugf("failed to listen %s: %v", laddr.String(), e)
|
||||
i++
|
||||
}
|
||||
return nil, ErrPort
|
||||
@@ -164,13 +167,13 @@ func (a *Agent) gatherCandidatesLocal(networkTypes []NetworkType) {
|
||||
var wg sync.WaitGroup
|
||||
defer wg.Wait()
|
||||
|
||||
localIPs := localInterfaces(networkTypes)
|
||||
localIPs := a.localInterfaces(networkTypes)
|
||||
wg.Add(len(localIPs) * len(supportedNetworks))
|
||||
for _, ip := range localIPs {
|
||||
for _, network := range supportedNetworks {
|
||||
go func(network string, ip net.IP) {
|
||||
defer wg.Done()
|
||||
conn, err := listenUDP(int(a.portmax), int(a.portmin), network, &net.UDPAddr{IP: ip, Port: 0})
|
||||
conn, err := a.listenUDP(int(a.portmax), int(a.portmin), network, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err != nil {
|
||||
a.log.Warnf("could not listen %s %s\n", network, ip)
|
||||
return
|
||||
@@ -234,13 +237,13 @@ func (a *Agent) gatherCandidatesSrflx(urls []*URL, networkTypes []NetworkType) {
|
||||
}
|
||||
|
||||
hostPort := fmt.Sprintf("%s:%d", url.Host, url.Port)
|
||||
serverAddr, err := net.ResolveUDPAddr(network, hostPort)
|
||||
serverAddr, err := a.net.ResolveUDPAddr(network, hostPort)
|
||||
if err != nil {
|
||||
a.log.Warnf("failed to resolve stun host: %s: %v", hostPort, err)
|
||||
continue
|
||||
}
|
||||
|
||||
conn, err := listenUDP(int(a.portmax), int(a.portmin), network, &net.UDPAddr{IP: nil, Port: 0})
|
||||
conn, err := a.listenUDP(int(a.portmax), int(a.portmin), network, &net.UDPAddr{IP: nil, Port: 0})
|
||||
if err != nil {
|
||||
a.log.Warnf("Failed to listen on %s for %s: %v\n", conn.LocalAddr().String(), serverAddr.String(), err)
|
||||
continue
|
||||
@@ -306,48 +309,58 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
|
||||
return ErrPasswordEmpty
|
||||
}
|
||||
|
||||
raddr, err := net.ResolveUDPAddr(network, fmt.Sprintf("%s:%d", url.Host, url.Port))
|
||||
locConn, err := a.net.ListenPacket(network, "0.0.0.0:0")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
c, err := net.DialUDP(network, nil, raddr)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
client, clientErr := turnc.New(turnc.Options{
|
||||
Conn: c,
|
||||
Username: url.Username,
|
||||
Password: url.Password,
|
||||
})
|
||||
if clientErr != nil {
|
||||
return clientErr
|
||||
}
|
||||
allocation, allocErr := client.Allocate()
|
||||
if allocErr != nil {
|
||||
return allocErr
|
||||
}
|
||||
|
||||
laddr := c.LocalAddr().(*net.UDPAddr)
|
||||
ip := allocation.Relayed().IP
|
||||
port := allocation.Relayed().Port
|
||||
client, err := turn.NewClient(&turn.ClientConfig{
|
||||
TURNServerAddr: fmt.Sprintf("%s:%d", url.Host, url.Port),
|
||||
Conn: locConn,
|
||||
Username: url.Username,
|
||||
Password: url.Password,
|
||||
LoggerFactory: a.loggerFactory,
|
||||
Net: a.net,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = client.Listen()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
relayConn, err := client.Allocate()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
laddr := locConn.LocalAddr().(*net.UDPAddr)
|
||||
raddr := relayConn.LocalAddr().(*net.UDPAddr)
|
||||
|
||||
relayConfig := CandidateRelayConfig{
|
||||
Network: network,
|
||||
Address: ip.String(),
|
||||
Port: port,
|
||||
Component: ComponentRTP,
|
||||
Address: raddr.IP.String(),
|
||||
Port: raddr.Port,
|
||||
RelAddr: laddr.IP.String(),
|
||||
RelPort: laddr.Port,
|
||||
OnClose: func() error {
|
||||
err2 := relayConn.Close()
|
||||
client.Close()
|
||||
return err2
|
||||
},
|
||||
}
|
||||
candidate, err := NewCandidateRelay(&relayConfig)
|
||||
if err != nil {
|
||||
a.log.Warnf("Failed to create server reflexive candidate: %s %s %d: %v\n", network, ip, port, err)
|
||||
a.log.Warnf("Failed to create relay candidate: %s %s: %v\n",
|
||||
network, raddr.String(), err)
|
||||
continue
|
||||
}
|
||||
candidate.setAllocation(client, allocation)
|
||||
|
||||
a.addCandidate(candidate)
|
||||
candidate.start(a, nil)
|
||||
candidate.start(a, relayConn)
|
||||
}
|
||||
|
||||
return nil
|
||||
@@ -357,7 +370,7 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
|
||||
// the XORMappedAddress returned by the stun server.
|
||||
//
|
||||
// Adapted from stun v0.2.
|
||||
func getXORMappedAddr(conn *net.UDPConn, serverAddr net.Addr, deadline time.Duration) (*stun.XORMappedAddress, error) {
|
||||
func getXORMappedAddr(conn net.PacketConn, serverAddr net.Addr, deadline time.Duration) (*stun.XORMappedAddress, error) {
|
||||
if deadline > 0 {
|
||||
if err := conn.SetReadDeadline(time.Now().Add(deadline)); err != nil {
|
||||
return nil, err
|
||||
@@ -369,7 +382,10 @@ func getXORMappedAddr(conn *net.UDPConn, serverAddr net.Addr, deadline time.Dura
|
||||
}
|
||||
}()
|
||||
resp, err := stunRequest(
|
||||
conn.Read,
|
||||
func(p []byte) (int, error) {
|
||||
n, _, errr := conn.ReadFrom(p)
|
||||
return n, errr
|
||||
},
|
||||
func(b []byte) (int, error) {
|
||||
return conn.WriteTo(b, serverAddr)
|
||||
},
|
||||
|
||||
+12
-4
@@ -6,26 +6,34 @@ import (
|
||||
)
|
||||
|
||||
func TestListenUDP(t *testing.T) {
|
||||
localIPs := localInterfaces([]NetworkType{NetworkTypeUDP4})
|
||||
a, err := NewAgent(&AgentConfig{})
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create agent: %s", err)
|
||||
}
|
||||
|
||||
localIPs := a.localInterfaces([]NetworkType{NetworkTypeUDP4})
|
||||
if len(localIPs) == 0 {
|
||||
t.Fatal("localInterfaces found no interfaces, unable to test")
|
||||
}
|
||||
|
||||
ip := localIPs[0]
|
||||
|
||||
conn, err := listenUDP(0, 0, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
conn, err := a.listenUDP(0, 0, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err != nil {
|
||||
t.Fatalf("listenUDP error with no port restriction %v", err)
|
||||
} else if conn == nil {
|
||||
t.Fatalf("listenUDP error with no port restriction return a nil conn")
|
||||
}
|
||||
|
||||
_, err = listenUDP(500, 499, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
_, err = a.listenUDP(4999, 5000, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err == nil {
|
||||
t.Fatal("listenUDP with invalid port range did not fail")
|
||||
}
|
||||
if err != ErrPort {
|
||||
t.Fatal("listenUDP with invalid port range did not return ErrPort")
|
||||
}
|
||||
|
||||
conn, err = listenUDP(5000, 5000, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
conn, err = a.listenUDP(5000, 5000, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err != nil {
|
||||
t.Fatalf("listenUDP error with no port restriction %v", err)
|
||||
} else if conn == nil {
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
package ice
|
||||
|
||||
import (
|
||||
"net"
|
||||
"testing"
|
||||
|
||||
"github.com/pion/logging"
|
||||
"github.com/pion/transport/vnet"
|
||||
)
|
||||
|
||||
func TestVNetGather(t *testing.T) {
|
||||
loggerFactory := logging.NewDefaultLoggerFactory()
|
||||
//log := loggerFactory.NewLogger("test")
|
||||
|
||||
t.Run("No local IP address", func(t *testing.T) {
|
||||
a, err := NewAgent(&AgentConfig{
|
||||
Net: vnet.NewNet(&vnet.NetConfig{}),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create agent: %s", err)
|
||||
}
|
||||
|
||||
localIPs := a.localInterfaces([]NetworkType{NetworkTypeUDP4})
|
||||
if len(localIPs) > 0 {
|
||||
t.Fatal("should return no local IP")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Gather a dynamic IP address", func(t *testing.T) {
|
||||
cider := "1.2.3.0/24"
|
||||
_, ipNet, err := net.ParseCIDR(cider)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to parse CIDR: %s", err)
|
||||
}
|
||||
|
||||
r, err := vnet.NewRouter(&vnet.RouterConfig{
|
||||
CIDR: cider,
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create a router: %s", err)
|
||||
}
|
||||
|
||||
nw := vnet.NewNet(&vnet.NetConfig{})
|
||||
if nw == nil {
|
||||
t.Fatalf("Failed to create a Net: %s", err)
|
||||
}
|
||||
|
||||
err = r.AddNet(nw)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to add a Net to the router: %s", err)
|
||||
}
|
||||
|
||||
a, err := NewAgent(&AgentConfig{
|
||||
Net: nw,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create agent: %s", err)
|
||||
}
|
||||
|
||||
localIPs := a.localInterfaces([]NetworkType{NetworkTypeUDP4})
|
||||
if len(localIPs) == 0 {
|
||||
t.Fatal("should have one local IP")
|
||||
}
|
||||
|
||||
for _, ip := range localIPs {
|
||||
if ip.IsLoopback() {
|
||||
t.Fatal("should not return loopback IP")
|
||||
}
|
||||
if !ipNet.Contains(ip) {
|
||||
t.Fatal("should be contained in the CIDR")
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("listenUDP", func(t *testing.T) {
|
||||
r, err := vnet.NewRouter(&vnet.RouterConfig{
|
||||
CIDR: "1.2.3.0/24",
|
||||
LoggerFactory: loggerFactory,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create a router: %s", err)
|
||||
}
|
||||
|
||||
nw := vnet.NewNet(&vnet.NetConfig{})
|
||||
if nw == nil {
|
||||
t.Fatalf("Failed to create a Net: %s", err)
|
||||
}
|
||||
|
||||
err = r.AddNet(nw)
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to add a Net to the router: %s", err)
|
||||
}
|
||||
|
||||
a, err := NewAgent(&AgentConfig{Net: nw})
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to create agent: %s", err)
|
||||
}
|
||||
|
||||
localIPs := a.localInterfaces([]NetworkType{NetworkTypeUDP4})
|
||||
if len(localIPs) == 0 {
|
||||
t.Fatal("localInterfaces found no interfaces, unable to test")
|
||||
}
|
||||
|
||||
ip := localIPs[0]
|
||||
|
||||
conn, err := a.listenUDP(0, 0, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err != nil {
|
||||
t.Fatalf("listenUDP error with no port restriction %v", err)
|
||||
} else if conn == nil {
|
||||
t.Fatalf("listenUDP error with no port restriction return a nil conn")
|
||||
}
|
||||
|
||||
_, err = a.listenUDP(4999, 5000, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err != ErrPort {
|
||||
t.Fatal("listenUDP with invalid port range did not return ErrPort")
|
||||
}
|
||||
|
||||
conn, err = a.listenUDP(5000, 5000, udp, &net.UDPAddr{IP: ip, Port: 0})
|
||||
if err != nil {
|
||||
t.Fatalf("listenUDP error with no port restriction %v", err)
|
||||
} else if conn == nil {
|
||||
t.Fatalf("listenUDP error with no port restriction return a nil conn")
|
||||
}
|
||||
|
||||
_, port, err := net.SplitHostPort(conn.LocalAddr().String())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
} else if port != "5000" {
|
||||
t.Fatalf("listenUDP with port restriction of 5000 listened on incorrect port (%s)", port)
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -3,12 +3,12 @@ module github.com/pion/ice
|
||||
go 1.12
|
||||
|
||||
require (
|
||||
github.com/pion/logging v0.2.1
|
||||
github.com/pion/logging v0.2.2
|
||||
github.com/pion/mdns v0.0.2
|
||||
github.com/pion/stun v0.3.1
|
||||
github.com/pion/transport v0.7.0
|
||||
github.com/pion/turn v1.1.4
|
||||
github.com/pion/turnc v0.0.6
|
||||
github.com/pion/transport v0.8.5-0.20190714194729-a583f529f67b
|
||||
github.com/pion/turn v1.2.2-0.20190714200549-4378deb04e7a
|
||||
github.com/stretchr/testify v1.3.0
|
||||
golang.org/x/net v0.0.0-20190619014844-b5b0513f8c1b
|
||||
golang.org/x/net v0.0.0-20190628185345-da137c7871d7
|
||||
golang.org/x/sys v0.0.0-20190712062909-fae7ac547cb7 // indirect
|
||||
)
|
||||
|
||||
@@ -1,21 +1,21 @@
|
||||
github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8=
|
||||
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/gortc/turn v0.7.1/go.mod h1:3FZ+LvCZKCKu6YYgwuYPqEi3FqCtdjfSFnFqVQNwfjk=
|
||||
github.com/gortc/turn v0.7.3 h1:CE72C79erbcsfa6L/QDhKztcl2kDq1UK20ImrJWDt/w=
|
||||
github.com/gortc/turn v0.7.3/go.mod h1:gvguwaGAFyv5/9KrcW9MkCgHALYD+e99mSM7pSCYYho=
|
||||
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
|
||||
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/gortc/turn v0.8.0 h1:WWQi1jkoPmc2E7qgUMcZleveKikT9Ksi3QGIl8ZtY3Q=
|
||||
github.com/gortc/turn v0.8.0/go.mod h1:gvguwaGAFyv5/9KrcW9MkCgHALYD+e99mSM7pSCYYho=
|
||||
github.com/pion/logging v0.2.1 h1:LwASkBKZ+2ysGJ+jLv1E/9H1ge0k1nTfi1X+5zirkDk=
|
||||
github.com/pion/logging v0.2.1/go.mod h1:k0/tDVsRCX2Mb2ZEmTqNa7CWsQPc+YYCB7Q+5pahoms=
|
||||
github.com/pion/logging v0.2.2 h1:M9+AIj/+pxNsDfAT64+MAVgJO0rsyLnoJKCqf//DoeY=
|
||||
github.com/pion/logging v0.2.2/go.mod h1:k0/tDVsRCX2Mb2ZEmTqNa7CWsQPc+YYCB7Q+5pahoms=
|
||||
github.com/pion/mdns v0.0.2 h1:T22Gg4dSuYVYsZ21oRFh9z7twzAm27+5PEKiABbjCvM=
|
||||
github.com/pion/mdns v0.0.2/go.mod h1:VrN3wefVgtfL8QgpEblPUC46ag1reLIfpqekCnKunLE=
|
||||
github.com/pion/stun v0.3.0/go.mod h1:xrCld6XM+6GWDZdvjPlLMsTU21rNxnO6UO8XsAvHr/M=
|
||||
github.com/pion/stun v0.3.1 h1:d09JJzOmOS8ZzIp8NppCMgrxGZpJ4Ix8qirfNYyI3BA=
|
||||
github.com/pion/stun v0.3.1/go.mod h1:xrCld6XM+6GWDZdvjPlLMsTU21rNxnO6UO8XsAvHr/M=
|
||||
github.com/pion/transport v0.7.0 h1:EsXN8TglHMlKZMo4ZGqwK6QgXBu0WYg7wfGMWIXsS+w=
|
||||
github.com/pion/transport v0.7.0/go.mod h1:iWZ07doqOosSLMhZ+FXUTq+TamDoXSllxpbGcfkCmbE=
|
||||
github.com/pion/turn v1.1.4 h1:yGxcasBvge4idNjxjowePn8oW43C4v70bXroBBKLyKY=
|
||||
github.com/pion/turn v1.1.4/go.mod h1:2O2GFDGO6+hJ5gsyExDhoNHtVcacPB1NOyc81gkq0WA=
|
||||
github.com/pion/turnc v0.0.6 h1:FHsmwYvdJ8mhT1/ZtWWer9L0unEb7AyRgrymfWy6mEY=
|
||||
github.com/pion/turnc v0.0.6/go.mod h1:4MSFv5i0v3MRkDLdo5eF9cD/xJtj1pxSphHNnxKL2W8=
|
||||
github.com/pion/transport v0.8.5-0.20190714194729-a583f529f67b h1:SIZcX3qOWle4QGeA85nNo0clrmcxQP/OqSYAFVPk/nw=
|
||||
github.com/pion/transport v0.8.5-0.20190714194729-a583f529f67b/go.mod h1:nAmRRnn+ArVtsoNuwktvAD+jrjSD7pA+H3iRmZwdUno=
|
||||
github.com/pion/turn v1.2.2-0.20190714200549-4378deb04e7a h1:kEgD9Rh5xjwJ9hU++Snd9i3pD1cVykUWNHbTbeDkpx0=
|
||||
github.com/pion/turn v1.2.2-0.20190714200549-4378deb04e7a/go.mod h1:x23bUN/QKPMch+Y0u5Jk1lx6AK4SHYCTRo384H56N3E=
|
||||
github.com/pkg/errors v0.8.1 h1:iURUrRGxPUNPdy5/HRSm+Yj6okJ6UtLINN0Q9M4+h3I=
|
||||
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
|
||||
@@ -28,6 +28,10 @@ golang.org/x/net v0.0.0-20190403144856-b630fd6fe46b h1:/zjbcJPEGAyu6Is/VBOALsgdi
|
||||
golang.org/x/net v0.0.0-20190403144856-b630fd6fe46b/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
|
||||
golang.org/x/net v0.0.0-20190619014844-b5b0513f8c1b h1:lkjdUzSyJ5P1+eal9fxXX9Xg2BTfswsonKUse48C0uE=
|
||||
golang.org/x/net v0.0.0-20190619014844-b5b0513f8c1b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
|
||||
golang.org/x/net v0.0.0-20190628185345-da137c7871d7 h1:rTIdg5QFRR7XCaK4LCjBiPbx8j4DQRpdYMnGn/bJUEU=
|
||||
golang.org/x/net v0.0.0-20190628185345-da137c7871d7/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
|
||||
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a h1:1BGLXjeY4akVXGgbC9HugT3Jv3hCI0z56oJR5vAMgBU=
|
||||
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20190712062909-fae7ac547cb7 h1:LepdCS8Gf/MVejFIt8lsiexZATdoGVyp5bcyS+rYoUI=
|
||||
golang.org/x/sys v0.0.0-20190712062909-fae7ac547cb7/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
|
||||
Reference in New Issue
Block a user