From 50bd9f60bd971b6d048fe89b882bb9bfd45a5582 Mon Sep 17 00:00:00 2001 From: Sean DuBois Date: Sat, 22 Feb 2020 22:21:19 -0800 Subject: [PATCH] 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 --- agent.go | 3 ++ gather.go | 140 ++++++++++++++++++++++++++++++++---------------------- 2 files changed, 87 insertions(+), 56 deletions(-) diff --git a/agent.go b/agent.go index 675f181..f1152ad 100644 --- a/agent.go +++ b/agent.go @@ -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 } } diff --git a/gather.go b/gather.go index fe06247..d7e35aa 100644 --- a/gather.go +++ b/gather.go @@ -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) } }