From d84503fa75ff4c59667667fc576e2459bc27ca6c Mon Sep 17 00:00:00 2001 From: giongto35 Date: Sun, 31 May 2020 03:28:36 +0800 Subject: [PATCH] Close channel to avoid routine leak (#188) --- pkg/webrtc/webrtc.go | 7 +++++++ pkg/worker/room/media.go | 10 ++++++++++ pkg/worker/room/room.go | 4 ++++ 3 files changed, 21 insertions(+) diff --git a/pkg/webrtc/webrtc.go b/pkg/webrtc/webrtc.go index 81b05616..877c3505 100644 --- a/pkg/webrtc/webrtc.go +++ b/pkg/webrtc/webrtc.go @@ -202,6 +202,11 @@ func (w *WebRTC) StartClient(isMobile bool, iceCB OnIceCallback) (string, error) log.Println("Received Voice from Client") for { + if w.RoomID == "" { + // skip sending voice when game is not running + continue + } + i, err := remoteTrack.Read(rtpBuf) // TODO: can receive track but the voice doesn't work if err == nil { @@ -290,6 +295,8 @@ func (w *WebRTC) StopClient() { // NOTE: ImageChannel is waiting for input. Close in writer is not correct for this close(w.ImageChannel) close(w.AudioChannel) + close(w.VoiceInChannel) + close(w.VoiceOutChannel) } // IsConnected comment diff --git a/pkg/worker/room/media.go b/pkg/worker/room/media.go index 203efc83..5718ed3c 100644 --- a/pkg/worker/room/media.go +++ b/pkg/worker/room/media.go @@ -49,6 +49,13 @@ func (r *Room) startVoice() { // broadcast voice go func() { for sample := range r.voiceInChannel { + r.voiceOutChannel <- sample + } + }() + + // fanout voice + go func() { + for sample := range r.voiceOutChannel { for _, webRTC := range r.rtcSessions { if webRTC.IsConnected() { // NOTE: can block here @@ -56,6 +63,9 @@ func (r *Room) startVoice() { } } } + for _, webRTC := range r.rtcSessions { + close(webRTC.VoiceOutChannel) + } }() } diff --git a/pkg/worker/room/room.go b/pkg/worker/room/room.go index ab35697f..4258e15c 100644 --- a/pkg/worker/room/room.go +++ b/pkg/worker/room/room.go @@ -222,6 +222,7 @@ func (r *Room) startWebRTCSession(peerconnection *webrtc.WebRTC) { log.Println("Start WebRTC session") go func() { + // set up voice input and output. A room has multiple voice input and only one combined voice output. for voiceInput := range peerconnection.VoiceInChannel { // NOTE: when room is no longer running. InputChannel needs to have extra event to go inside the loop @@ -262,6 +263,7 @@ func (r *Room) RemoveSession(w *webrtc.WebRTC) { log.Println("found session: ", w.ID) if s.ID == w.ID { r.rtcSessions = append(r.rtcSessions[:i], r.rtcSessions[i+1:]...) + s.RoomID = "" log.Println("Removed session ", s.ID, " from room: ", r.ID) break } @@ -309,6 +311,8 @@ func (r *Room) Close() { } log.Println("Closing input of room ", r.ID) close(r.inputChannel) + close(r.voiceOutChannel) + close(r.voiceInChannel) close(r.Done) // Close here is a bit wrong because this read channel // Just dont close it, let it be gc