Browse Source

webrtc: fix race condition that caused random crashes during handshake (#2072)

pull/2073/head
Alessandro Ros 3 years ago committed by GitHub
parent
commit
5066ba403c
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
  1. 14
      internal/core/webrtc_manager.go
  2. 59
      internal/core/webrtc_session.go

14
internal/core/webrtc_manager.go

@ -269,7 +269,7 @@ outer:
sx := newWebRTCSession( sx := newWebRTCSession(
m.ctx, m.ctx,
m.readBufferCount, m.readBufferCount,
req, req.remoteAddr,
&wg, &wg,
m.iceHostNAT1To1IPs, m.iceHostNAT1To1IPs,
m.iceUDPMux, m.iceUDPMux,
@ -385,15 +385,9 @@ func (m *webRTCManager) sessionNew(req webRTCSessionNewReq) webRTCSessionNewRes
select { select {
case m.chSessionNew <- req: case m.chSessionNew <- req:
res1 := <-req.res res := <-req.res
select {
case res2 := <-req.res:
return res2
case <-res1.sx.ctx.Done(): return res.sx.new(req)
return webRTCSessionNewRes{err: fmt.Errorf("terminated"), errStatusCode: http.StatusInternalServerError}
}
case <-m.ctx.Done(): case <-m.ctx.Done():
return webRTCSessionNewRes{err: fmt.Errorf("terminated"), errStatusCode: http.StatusInternalServerError} return webRTCSessionNewRes{err: fmt.Errorf("terminated"), errStatusCode: http.StatusInternalServerError}
@ -420,7 +414,7 @@ func (m *webRTCManager) sessionAddCandidates(
return res1 return res1
} }
return res1.sx.addRemoteCandidates(req) return res1.sx.addCandidates(req)
case <-m.ctx.Done(): case <-m.ctx.Done():
return webRTCSessionAddCandidatesRes{err: fmt.Errorf("terminated")} return webRTCSessionAddCandidatesRes{err: fmt.Errorf("terminated")}

59
internal/core/webrtc_session.go

@ -113,7 +113,6 @@ type webRTCSessionPathManager interface {
type webRTCSession struct { type webRTCSession struct {
readBufferCount int readBufferCount int
req webRTCSessionNewReq
wg *sync.WaitGroup wg *sync.WaitGroup
iceHostNAT1To1IPs []string iceHostNAT1To1IPs []string
iceUDPMux ice.UDPMux iceUDPMux ice.UDPMux
@ -126,17 +125,19 @@ type webRTCSession struct {
created time.Time created time.Time
uuid uuid.UUID uuid uuid.UUID
secret uuid.UUID secret uuid.UUID
req webRTCSessionNewReq
answerSent bool answerSent bool
mutex sync.RWMutex mutex sync.RWMutex
pc *peerConnection pc *peerConnection
chAddRemoteCandidates chan webRTCSessionAddCandidatesReq chNew chan webRTCSessionNewReq
chAddCandidates chan webRTCSessionAddCandidatesReq
} }
func newWebRTCSession( func newWebRTCSession(
parentCtx context.Context, parentCtx context.Context,
readBufferCount int, readBufferCount int,
req webRTCSessionNewReq, remoteAddr string,
wg *sync.WaitGroup, wg *sync.WaitGroup,
iceHostNAT1To1IPs []string, iceHostNAT1To1IPs []string,
iceUDPMux ice.UDPMux, iceUDPMux ice.UDPMux,
@ -148,7 +149,6 @@ func newWebRTCSession(
s := &webRTCSession{ s := &webRTCSession{
readBufferCount: readBufferCount, readBufferCount: readBufferCount,
req: req,
wg: wg, wg: wg,
iceHostNAT1To1IPs: iceHostNAT1To1IPs, iceHostNAT1To1IPs: iceHostNAT1To1IPs,
iceUDPMux: iceUDPMux, iceUDPMux: iceUDPMux,
@ -160,10 +160,11 @@ func newWebRTCSession(
created: time.Now(), created: time.Now(),
uuid: uuid.New(), uuid: uuid.New(),
secret: uuid.New(), secret: uuid.New(),
chAddRemoteCandidates: make(chan webRTCSessionAddCandidatesReq), chNew: make(chan webRTCSessionNewReq),
chAddCandidates: make(chan webRTCSessionAddCandidatesReq),
} }
s.Log(logger.Info, "created by %s", req.remoteAddr) s.Log(logger.Info, "created by %s", remoteAddr)
wg.Add(1) wg.Add(1)
go s.run() go s.run()
@ -183,7 +184,25 @@ func (s *webRTCSession) close() {
func (s *webRTCSession) run() { func (s *webRTCSession) run() {
defer s.wg.Done() defer s.wg.Done()
errStatusCode, err := s.runInner() err := s.runInner()
s.ctxCancel()
s.parent.sessionClose(s)
s.Log(logger.Info, "closed (%v)", err)
}
func (s *webRTCSession) runInner() error {
select {
case req := <-s.chNew:
s.req = req
case <-s.ctx.Done():
return fmt.Errorf("terminated")
}
errStatusCode, err := s.runInner2()
if !s.answerSent { if !s.answerSent {
select { select {
@ -195,14 +214,10 @@ func (s *webRTCSession) run() {
} }
} }
s.ctxCancel() return err
s.parent.sessionClose(s)
s.Log(logger.Info, "closed (%v)", err)
} }
func (s *webRTCSession) runInner() (int, error) { func (s *webRTCSession) runInner2() (int, error) {
if s.req.publish { if s.req.publish {
return s.runPublish() return s.runPublish()
} }
@ -495,7 +510,7 @@ outer:
func (s *webRTCSession) readRemoteCandidates(pc *peerConnection) { func (s *webRTCSession) readRemoteCandidates(pc *peerConnection) {
for { for {
select { select {
case req := <-s.chAddRemoteCandidates: case req := <-s.chAddCandidates:
for _, candidate := range req.candidates { for _, candidate := range req.candidates {
err := pc.AddICECandidate(*candidate) err := pc.AddICECandidate(*candidate)
if err != nil { if err != nil {
@ -510,11 +525,23 @@ func (s *webRTCSession) readRemoteCandidates(pc *peerConnection) {
} }
} }
func (s *webRTCSession) addRemoteCandidates( // new is called by webRTCHTTPServer through webRTCManager.
func (s *webRTCSession) new(req webRTCSessionNewReq) webRTCSessionNewRes {
select {
case s.chNew <- req:
return <-req.res
case <-s.ctx.Done():
return webRTCSessionNewRes{err: fmt.Errorf("terminated"), errStatusCode: http.StatusInternalServerError}
}
}
// addCandidates is called by webRTCHTTPServer through webRTCManager.
func (s *webRTCSession) addCandidates(
req webRTCSessionAddCandidatesReq, req webRTCSessionAddCandidatesReq,
) webRTCSessionAddCandidatesRes { ) webRTCSessionAddCandidatesRes {
select { select {
case s.chAddRemoteCandidates <- req: case s.chAddCandidates <- req:
return <-req.res return <-req.res
case <-s.ctx.Done(): case <-s.ctx.Done():

Loading…
Cancel
Save