chore: Update pion/ice fork to resolve goroutine leak (#78)

* chore: Update pion/ice fork to resolve goroutine leak

* Flush remote too

* Add logs for setting the description

* Try locking only on remote

* Remove local bufferring in favor of remote

* Remove unused flush func

* Set candidates flushed to true

* Defer flush until the end of negotiation

* Buffer ICE candidates

* Add comment clarifying channel buffer

* Flush after handshake

* Move away from fork

* Ignore pion/ice leaks
This commit is contained in:
Kyle Carberry
2022-01-27 20:43:55 -06:00
committed by GitHub
parent 30dae97c3e
commit 27f7299383
4 changed files with 46 additions and 47 deletions
+33 -41
View File
@@ -62,17 +62,19 @@ func newWithClientOrServer(servers []webrtc.ICEServer, client bool, opts *ConnOp
return nil, xerrors.Errorf("create peer connection: %w", err)
}
conn := &Conn{
pingChannelID: 1,
pingEchoChannelID: 2,
opts: opts,
rtc: rtc,
offerrer: client,
closed: make(chan struct{}),
dcOpenChannel: make(chan *webrtc.DataChannel),
dcDisconnectChannel: make(chan struct{}),
dcFailedChannel: make(chan struct{}),
localCandidateChannel: make(chan webrtc.ICECandidateInit),
pendingCandidates: make([]webrtc.ICECandidateInit, 0),
pingChannelID: 1,
pingEchoChannelID: 2,
opts: opts,
rtc: rtc,
offerrer: client,
closed: make(chan struct{}),
dcOpenChannel: make(chan *webrtc.DataChannel),
dcDisconnectChannel: make(chan struct{}),
dcFailedChannel: make(chan struct{}),
// This channel needs to be bufferred otherwise slow consumers
// of this will cause a connection failure.
localCandidateChannel: make(chan webrtc.ICECandidateInit, 16),
pendingRemoteCandidates: make([]webrtc.ICECandidateInit, 0),
localSessionDescriptionChannel: make(chan webrtc.SessionDescription),
remoteSessionDescriptionChannel: make(chan webrtc.SessionDescription),
}
@@ -120,7 +122,7 @@ type Conn struct {
localSessionDescriptionChannel chan webrtc.SessionDescription
remoteSessionDescriptionChannel chan webrtc.SessionDescription
pendingCandidates []webrtc.ICECandidateInit
pendingRemoteCandidates []webrtc.ICECandidateInit
pendingCandidatesMutex sync.Mutex
pendingCandidatesFlushed bool
@@ -142,14 +144,6 @@ func (c *Conn) init() error {
if iceCandidate == nil {
return
}
c.pendingCandidatesMutex.Lock()
defer c.pendingCandidatesMutex.Unlock()
if !c.pendingCandidatesFlushed {
c.opts.Logger.Debug(context.Background(), "adding local candidate to buffer")
c.pendingCandidates = append(c.pendingCandidates, iceCandidate.ToJSON())
return
}
c.opts.Logger.Debug(context.Background(), "adding local candidate")
select {
case <-c.closed:
@@ -262,6 +256,7 @@ func (c *Conn) negotiate() {
_ = c.CloseWithError(xerrors.Errorf("create offer: %w", err))
return
}
c.opts.Logger.Debug(context.Background(), "setting local description")
err = c.rtc.SetLocalDescription(offer)
if err != nil {
_ = c.CloseWithError(xerrors.Errorf("set local description: %w", err))
@@ -281,25 +276,20 @@ func (c *Conn) negotiate() {
case remoteDescription = <-c.remoteSessionDescriptionChannel:
}
c.opts.Logger.Debug(context.Background(), "setting remote description")
err := c.rtc.SetRemoteDescription(remoteDescription)
if err != nil {
c.pendingCandidatesMutex.Unlock()
_ = c.CloseWithError(xerrors.Errorf("set remote description (closed %v): %w", c.isClosed(), err))
return
}
if c.offerrer {
// ICE candidates reset when an offer/answer is set for the first
// time. If candidates flush before this point, a connection could fail.
c.flushPendingCandidates()
}
if !c.offerrer {
answer, err := c.rtc.CreateAnswer(&webrtc.AnswerOptions{})
if err != nil {
_ = c.CloseWithError(xerrors.Errorf("create answer: %w", err))
return
}
c.opts.Logger.Debug(context.Background(), "setting local description")
err = c.rtc.SetLocalDescription(answer)
if err != nil {
_ = c.CloseWithError(xerrors.Errorf("set local description: %w", err))
@@ -313,28 +303,23 @@ func (c *Conn) negotiate() {
return
case c.localSessionDescriptionChannel <- answer:
}
// Wait until the local description is set to flush candidates.
c.flushPendingCandidates()
}
}
// flushPendingCandidates writes all local candidates to the candidate send channel.
// The localCandidateChannel is expected to be serviced, otherwise this could block.
func (c *Conn) flushPendingCandidates() {
// The ICE transport resets when the remote description is updated.
// Adding ICE candidates before this point causes a failed connection,
// because the candidate would be lost.
c.pendingCandidatesMutex.Lock()
defer c.pendingCandidatesMutex.Unlock()
for _, pendingCandidate := range c.pendingCandidates {
c.opts.Logger.Debug(context.Background(), "flushing local candidate")
select {
case <-c.closed:
for _, pendingCandidate := range c.pendingRemoteCandidates {
c.opts.Logger.Debug(context.Background(), "flushing remote candidate")
err := c.rtc.AddICECandidate(pendingCandidate)
if err != nil {
_ = c.CloseWithError(xerrors.Errorf("flush pending candidates: %w", err))
return
case c.localCandidateChannel <- pendingCandidate:
}
}
c.pendingCandidates = make([]webrtc.ICECandidateInit, 0)
c.pendingCandidatesFlushed = true
c.opts.Logger.Debug(context.Background(), "flushed candidates")
c.opts.Logger.Debug(context.Background(), "flushed remote candidates")
}
// LocalCandidate returns a channel that emits when a local candidate
@@ -345,6 +330,13 @@ func (c *Conn) LocalCandidate() <-chan webrtc.ICECandidateInit {
// AddRemoteCandidate adds a remote candidate to the RTC connection.
func (c *Conn) AddRemoteCandidate(i webrtc.ICECandidateInit) error {
c.pendingCandidatesMutex.Lock()
defer c.pendingCandidatesMutex.Unlock()
if !c.pendingCandidatesFlushed {
c.opts.Logger.Debug(context.Background(), "adding remote candidate to buffer")
c.pendingRemoteCandidates = append(c.pendingRemoteCandidates, i)
return nil
}
c.opts.Logger.Debug(context.Background(), "adding remote candidate")
return c.rtc.AddICECandidate(i)
}
+8 -1
View File
@@ -48,7 +48,14 @@ var (
)
func TestMain(m *testing.M) {
goleak.VerifyTestMain(m)
// pion/ice doesn't properly close immediately. The solution for this isn't yet known. See:
// https://github.com/pion/ice/pull/413
goleak.VerifyTestMain(m,
goleak.IgnoreTopFunction("github.com/pion/ice/v2.(*Agent).startOnConnectionStateChangeRoutine.func1"),
goleak.IgnoreTopFunction("github.com/pion/ice/v2.(*Agent).startOnConnectionStateChangeRoutine.func2"),
goleak.IgnoreTopFunction("github.com/pion/ice/v2.(*Agent).taskLoop"),
goleak.IgnoreTopFunction("internal/poll.runtime_pollWait"),
)
}
func TestConn(t *testing.T) {