Close sockets if error during gathering

Add logic that if we hit errors close the socket
we are using that hasn't been moved into a a candidate yet.

Resolves #102
This commit is contained in:
Sean DuBois
2020-02-23 22:05:13 -08:00
parent 33bb85884f
commit 50bd9f60bd
2 changed files with 87 additions and 56 deletions
+3
View File
@@ -891,6 +891,9 @@ func (a *Agent) addCandidate(c Candidate, candidateConn net.PacketConn) error {
set := a.localCandidates[c.NetworkType()]
for _, candidate := range set {
if candidate.Equal(c) {
if err := c.close(); err != nil {
a.log.Warnf("Failed to close duplicate candidate: %v", err)
}
return
}
}
+84 -56
View File
@@ -3,9 +3,9 @@ package ice
import (
"fmt"
"net"
"sync"
"time"
"github.com/pion/logging"
"github.com/pion/turn/v2"
)
@@ -13,6 +13,18 @@ const (
stunGatherTimeout = time.Second * 5
)
type closeable interface {
Close() error
}
// Close a net.Conn and log if we have a failure
func closeConnAndLog(c closeable, log logging.LeveledLogger, msg string) {
log.Warnf(msg)
if err := c.Close(); err != nil {
log.Warnf("Failed to close conn: %v", err)
}
}
// GatherCandidates initiates the trickle based gathering process.
func (a *Agent) GatherCandidates() error {
gatherErrChan := make(chan error, 1)
@@ -81,16 +93,12 @@ func (a *Agent) gatherCandidates() {
}
func (a *Agent) gatherCandidatesLocal(networkTypes []NetworkType) {
var wg sync.WaitGroup
defer wg.Wait()
localIPs, err := localInterfaces(a.net, a.interfaceFilter, networkTypes)
if err != nil {
a.log.Warnf("failed to iterate local interfaces, host candidates will not be gathered %s", err)
return
}
wg.Add(len(localIPs) * len(supportedNetworks))
for _, ip := range localIPs {
mappedIP := ip
if a.mDNSMode != MulticastDNSModeQueryAndGather && a.extIPMapper != nil && a.extIPMapper.candidateType == CandidateTypeHost {
@@ -101,46 +109,45 @@ func (a *Agent) gatherCandidatesLocal(networkTypes []NetworkType) {
}
}
address := mappedIP.String()
if a.mDNSMode == MulticastDNSModeQueryAndGather {
address = a.mDNSName
}
for _, network := range supportedNetworks {
go func(network string, ip, mappedIP net.IP) {
defer wg.Done()
conn, err := listenUDPInPortRange(a.net, a.log, 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
}
conn, err := listenUDPInPortRange(a.net, a.log, 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)
continue
}
address := mappedIP.String()
if a.mDNSMode == MulticastDNSModeQueryAndGather {
address = a.mDNSName
}
port := conn.LocalAddr().(*net.UDPAddr).Port
hostConfig := CandidateHostConfig{
Network: network,
Address: address,
Port: port,
Component: ComponentRTP,
}
port := conn.LocalAddr().(*net.UDPAddr).Port
c, err := NewCandidateHost(&hostConfig)
if err != nil {
closeConnAndLog(conn, a.log, fmt.Sprintf("Failed to create host candidate: %s %s %d: %v\n", network, mappedIP, port, err))
continue
}
hostConfig := CandidateHostConfig{
Network: network,
Address: address,
Port: port,
Component: ComponentRTP,
if a.mDNSMode == MulticastDNSModeQueryAndGather {
if err = c.setIP(ip); err != nil {
closeConnAndLog(conn, a.log, fmt.Sprintf("Failed to create host candidate: %s %s %d: %v\n", network, mappedIP, port, err))
continue
}
}
c, err := NewCandidateHost(&hostConfig)
if err != nil {
a.log.Warnf("Failed to create host candidate: %s %s %d: %v\n", network, mappedIP, port, err)
return
if err := a.addCandidate(c, conn); err != nil {
if closeErr := c.close(); closeErr != nil {
a.log.Warnf("Failed to close candidate: %v", closeErr)
}
if a.mDNSMode == MulticastDNSModeQueryAndGather {
if err = c.setIP(ip); err != nil {
a.log.Warnf("Failed to create host candidate: %s %s %d: %v\n", network, mappedIP, port, err)
return
}
}
if err := a.addCandidate(c, conn); err != nil {
a.log.Warnf("Failed to append to localCandidates and run onCandidateHdlr: %v\n", err)
}
}(network, ip, mappedIP)
a.log.Warnf("Failed to append to localCandidates and run onCandidateHdlr: %v\n", err)
}
}
}
}
@@ -158,7 +165,7 @@ func (a *Agent) gatherCandidatesSrflxMapped(networkTypes []NetworkType) {
laddr := conn.LocalAddr().(*net.UDPAddr)
mappedIP, err := a.extIPMapper.findExternalIP(laddr.IP.String())
if err != nil {
a.log.Warnf("1:1 NAT mapping is enabled but no external IP is found for %s\n", laddr.IP.String())
closeConnAndLog(conn, a.log, fmt.Sprintf("1:1 NAT mapping is enabled but no external IP is found for %s\n", laddr.IP.String()))
continue
}
@@ -172,15 +179,18 @@ func (a *Agent) gatherCandidatesSrflxMapped(networkTypes []NetworkType) {
}
c, err := NewCandidateServerReflexive(&srflxConfig)
if err != nil {
a.log.Warnf("Failed to create server reflexive candidate: %s %s %d: %v\n",
closeConnAndLog(conn, a.log, fmt.Sprintf("Failed to create server reflexive candidate: %s %s %d: %v\n",
network,
mappedIP.String(),
laddr.Port,
err)
err))
continue
}
if err := a.addCandidate(c, conn); err != nil {
if closeErr := c.close(); closeErr != nil {
a.log.Warnf("Failed to close candidate: %v", closeErr)
}
a.log.Warnf("Failed to append to localCandidates and run onCandidateHdlr: %v\n", err)
}
}
@@ -210,13 +220,13 @@ func (a *Agent) gatherCandidatesSrflx(urls []*URL, networkTypes []NetworkType) {
conn, err := listenUDPInPortRange(a.net, a.log, int(a.portmax), int(a.portmin), network, &net.UDPAddr{IP: nil, Port: 0})
if err != nil {
a.log.Warnf("Failed to listen for %s: %v\n", serverAddr.String(), err)
closeConnAndLog(conn, a.log, fmt.Sprintf("Failed to listen for %s: %v\n", serverAddr.String(), err))
continue
}
xoraddr, err := getXORMappedAddr(conn, serverAddr, stunGatherTimeout)
if err != nil {
a.log.Warnf("could not get server reflexive address %s %s: %v\n", network, url, err)
closeConnAndLog(conn, a.log, fmt.Sprintf("could not get server reflexive address %s %s: %v\n", network, url, err))
continue
}
@@ -234,11 +244,14 @@ func (a *Agent) gatherCandidatesSrflx(urls []*URL, networkTypes []NetworkType) {
}
c, err := NewCandidateServerReflexive(&srflxConfig)
if err != nil {
a.log.Warnf("Failed to create server reflexive candidate: %s %s %d: %v\n", network, ip, port, err)
closeConnAndLog(conn, a.log, fmt.Sprintf("Failed to create server reflexive candidate: %s %s %d: %v\n", network, ip, port, err))
continue
}
if err := a.addCandidate(c, conn); err != nil {
if closeErr := c.close(); closeErr != nil {
a.log.Warnf("Failed to close candidate: %v", closeErr)
}
a.log.Warnf("Failed to append to localCandidates and run onCandidateHdlr: %v\n", err)
}
}
@@ -246,7 +259,6 @@ func (a *Agent) gatherCandidatesSrflx(urls []*URL, networkTypes []NetworkType) {
}
func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
network := NetworkTypeUDP4.String() // TODO IPv6
for _, url := range urls {
switch {
case url.Scheme != SchemeTypeTURN:
@@ -256,7 +268,10 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
case url.Password == "":
return ErrPasswordEmpty
}
}
network := NetworkTypeUDP4.String() // TODO IPv6
for _, url := range urls {
TURNServerAddr := fmt.Sprintf("%s:%d", url.Host, url.Port)
var (
locConn net.PacketConn
@@ -268,7 +283,8 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
if url.Proto == ProtoTypeUDP {
locConn, err = a.net.ListenPacket(network, "0.0.0.0:0")
if err != nil {
return err
a.log.Warnf("Failed to listen %s: %v\n", network, err)
continue
}
RelAddr = locConn.LocalAddr().(*net.UDPAddr).IP.String()
@@ -281,12 +297,14 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
tcpAddr, err = net.ResolveTCPAddr(NetworkTypeTCP4.String(), TURNServerAddr)
if err != nil {
return err
a.log.Warnf("Failed to resolve TCP Addr %s: %v\n", TURNServerAddr, err)
continue
}
tcpConn, err = net.DialTCP(NetworkTypeTCP4.String(), nil, tcpAddr)
if err != nil {
return err
a.log.Warnf("Failed to Dial TCP Addr %s: %v\n", TURNServerAddr, err)
continue
}
RelAddr = tcpConn.LocalAddr().(*net.TCPAddr).IP.String()
@@ -303,21 +321,24 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
Net: a.net,
})
if err != nil {
return err
closeConnAndLog(locConn, a.log, fmt.Sprintf("Failed to build new turn.Client %s %s\n", TURNServerAddr, err))
continue
}
err = client.Listen()
if err != nil {
return err
if err = client.Listen(); err != nil {
client.Close()
closeConnAndLog(locConn, a.log, fmt.Sprintf("Failed to listen on turn.Client %s %s\n", TURNServerAddr, err))
continue
}
relayConn, err := client.Allocate()
if err != nil {
return err
client.Close()
closeConnAndLog(locConn, a.log, fmt.Sprintf("Failed to allocate on turn.Client %s %s\n", TURNServerAddr, err))
continue
}
raddr := relayConn.LocalAddr().(*net.UDPAddr)
relayConfig := CandidateRelayConfig{
Network: network,
Component: ComponentRTP,
@@ -332,12 +353,19 @@ func (a *Agent) gatherCandidatesRelay(urls []*URL) error {
}
candidate, err := NewCandidateRelay(&relayConfig)
if err != nil {
a.log.Warnf("Failed to create relay candidate: %s %s: %v\n",
network, raddr.String(), err)
if relayConErr := relayConn.Close(); relayConErr != nil {
a.log.Warnf("Failed to close relay %v", relayConErr)
}
client.Close()
closeConnAndLog(locConn, a.log, fmt.Sprintf("Failed to create relay candidate: %s %s: %v\n", network, raddr.String(), err))
continue
}
if err := a.addCandidate(candidate, relayConn); err != nil {
if closeErr := candidate.close(); closeErr != nil {
a.log.Warnf("Failed to close candidate: %v", closeErr)
}
a.log.Warnf("Failed to append to localCandidates and run onCandidateHdlr: %v\n", err)
}
}