From 322b5bffe0c6df7d4becae57bddbad488a3c2cf4 Mon Sep 17 00:00:00 2001 From: giongto35 Date: Tue, 7 May 2019 02:26:52 +0800 Subject: [PATCH] Free up resource after room is closed --- emulator/director.go | 3 ++- handler/browser.go | 3 ++- handler/handlers.go | 10 ++++------ handler/overlord.go | 4 ++-- handler/room/media.go | 43 +++++++++++++++++++++++++++++-------------- handler/room/room.go | 23 +++++++++++++++++------ webrtc/webrtc.go | 5 ----- 7 files changed, 56 insertions(+), 35 deletions(-) diff --git a/emulator/director.go b/emulator/director.go index 142b45a0..a9f6e135 100644 --- a/emulator/director.go +++ b/emulator/director.go @@ -64,6 +64,7 @@ func (d *Director) Start(paths []string) { // audio := NewAudio() // audio.Start() // d.audio = audio + log.Println("Start game: ", paths) if len(paths) == 1 { d.PlayGame(paths[0]) @@ -95,7 +96,7 @@ L: func (d *Director) PlayGame(path string) { console, err := nes.NewConsole(path) if err != nil { - log.Fatalln(err) + log.Println("Err: Cannot load path, Got:", err) } // Set GameView as current view d.SetView(NewGameView(console, path, d.roomID, d.imageChannel, d.audioChannel, d.inputChannel)) diff --git a/handler/browser.go b/handler/browser.go index 98c81672..6c174f5f 100644 --- a/handler/browser.go +++ b/handler/browser.go @@ -14,7 +14,8 @@ type BrowserClient struct { *cws.Client } -func (s *Session) RegisterBrowserClient() { +// RouteBrowser are all routes server received from browser +func (s *Session) RouteBrowser() { browserClient := s.BrowserClient browserClient.Receive("heartbeat", func(resp cws.WSPacket) cws.WSPacket { diff --git a/handler/handlers.go b/handler/handlers.go index 70424ddb..c1b11ca0 100644 --- a/handler/handlers.go +++ b/handler/handlers.go @@ -1,7 +1,6 @@ package handler import ( - "fmt" "io/ioutil" "log" "net/http" @@ -84,8 +83,6 @@ func (h *Handler) WS(w http.ResponseWriter, r *http.Request) { } client := NewBrowserClient(c) - //client := NewClient(c) - ////sessionID := strconv.Itoa(rand.Int()) sessionID := uuid.Must(uuid.NewV4()).String() wssession := &Session{ ID: sessionID, @@ -94,14 +91,15 @@ func (h *Handler) WS(w http.ResponseWriter, r *http.Request) { peerconnection: webrtc.NewWebRTC(), handler: h, } + if wssession.OverlordClient != nil { - wssession.RegisterOverlordClient() + wssession.RouteOverlord() go wssession.OverlordClient.Heartbeat() go wssession.OverlordClient.Listen() } - wssession.RegisterBrowserClient() - fmt.Println("oclient : ", h.oClient) + wssession.RouteBrowser() + log.Println("oclient : ", h.oClient) wssession.BrowserClient.Send(cws.WSPacket{ ID: "gamelist", diff --git a/handler/overlord.go b/handler/overlord.go index a6a2c9df..76b18ce7 100644 --- a/handler/overlord.go +++ b/handler/overlord.go @@ -29,8 +29,8 @@ func NewOverlordClient(oc *websocket.Conn) *OverlordClient { return oclient } -// RegisterOverlordClient routes overlord Client -func (s *Session) RegisterOverlordClient() { +// RouteOverlord are all routes server received from overlord +func (s *Session) RouteOverlord() { oclient := s.OverlordClient // Received from overlord the serverID diff --git a/handler/room/media.go b/handler/room/media.go index 3a3581ab..aff2dbb9 100644 --- a/handler/room/media.go +++ b/handler/room/media.go @@ -28,9 +28,16 @@ func (r *Room) startAudio() { for { select { case <-r.Done: - r.Close() + log.Println("Room ", r.ID, " audio channel closed") return - case sample := <-r.audioChannel: + case sample, ok := <-r.audioChannel: + if !ok { + // Just for guarding + log.Println("Warn: Room ", r.ID, " audio channel closed unexpectedly") + return + } + + // TODO: Use worker pool for encoding pcm[idx] = sample idx++ if idx == len(pcm) { @@ -45,12 +52,13 @@ func (r *Room) startAudio() { data = data[:n] data = append(data, count) - r.sessionsLock.Lock() + // TODO: r.rtcSessions is rarely updated. Lock will hold down perf + //r.sessionsLock.Lock() for _, webRTC := range r.rtcSessions { // Client stopped - if webRTC.IsClosed() { - continue - } + //if !webRTC.IsClosed() { + //continue + //} // encode frame // fanout audioChannel @@ -60,7 +68,7 @@ func (r *Room) startAudio() { } //isRoomRunning = true } - r.sessionsLock.Unlock() + //r.sessionsLock.Unlock() idx = 0 count = (count + 1) & 0xff @@ -75,17 +83,24 @@ func (r *Room) startVideo() { for { select { case <-r.Done: - r.Close() + log.Println("Room ", r.ID, " video channel closed") return - case image := <-r.imageChannel: + case image, ok := <-r.imageChannel: + if !ok { + // Just for guarding, should not reached + log.Println("Warn: Room ", r.ID, " video channel closed unexpectedly") + return + } + // TODO: Use worker pool for encoding yuv := util.RgbaToYuv(image) - r.sessionsLock.Lock() + // TODO: r.rtcSessions is rarely updated. Lock will hold down perf + //r.sessionsLock.Lock() for _, webRTC := range r.rtcSessions { // Client stopped - if webRTC.IsClosed() { - continue - } + //if webRTC.IsClosed() { + //continue + //} // encode frame // fanout imageChannel @@ -94,7 +109,7 @@ func (r *Room) startVideo() { webRTC.ImageChannel <- yuv } } - r.sessionsLock.Unlock() + //r.sessionsLock.Unlock() } } } diff --git a/handler/room/room.go b/handler/room/room.go index a02277e3..a3ea2e10 100644 --- a/handler/room/room.go +++ b/handler/room/room.go @@ -65,6 +65,7 @@ func NewRoom(roomID, gamepath, gameName string, onlineStorage *storage.Client) * go room.startVideo() go room.startAudio() + // Check if room is on local storage, if not, pull from GCS to local storage path := gamepath + "/" + gameName go func(path, roomID string) { // Check room is on local or fetch from server @@ -109,12 +110,16 @@ func (r *Room) startWebRTCSession(peerconnection *webrtc.WebRTC, playerIndex int inputChannel := r.inputChannel for { select { + case <-r.Done: + log.Println("Close listening from peerconnection for room", r.ID) + return case <-peerconnection.Done: r.removeSession(peerconnection) default: } // Client stopped - if peerconnection.IsClosed() { + if !peerconnection.IsConnected() { + log.Println("peerconnection is closed", peerconnection) return } @@ -150,9 +155,12 @@ func (r *Room) removeSession(w *webrtc.WebRTC) { fmt.Println("found session: ", len(r.rtcSessions)) // If room has no sessions, close room + // Note: this logic cannot be brought outside of forloop because we only close room if room had at least one session if len(r.rtcSessions) == 0 { log.Println("No session in room") - r.Done <- struct{}{} + r.Close() + // can consider sanding close to room and room do clean + //close(r.Done) } break } @@ -162,9 +170,12 @@ func (r *Room) removeSession(w *webrtc.WebRTC) { func (r *Room) Close() { log.Println("Closing room", r) r.director.Done <- struct{}{} + close(r.Done) + // Not close fan input channel here, close in writer //close(r.inputChannel) - //close(r.imageChannel) - //close(r.audioChannel) + // Close fan output channel + close(r.imageChannel) + close(r.audioChannel) } func (r *Room) SaveGame() error { @@ -186,7 +197,7 @@ func (r *Room) SaveGame() error { } func (r *Room) loadRoomOnline(roomID string, savepath string) error { - log.Println("Loading game from cloud storage") + log.Println("Try loading game from cloud storage") // If the game is not on local server // Try to load from gcloud data, err := r.onlineStorage.LoadFile(roomID) @@ -212,7 +223,7 @@ func (r *Room) LoadGame() error { func (r *Room) IsRunning() bool { // If there is running session for _, s := range r.rtcSessions { - if !s.IsClosed() { + if s.IsConnected() { return true } } diff --git a/webrtc/webrtc.go b/webrtc/webrtc.go index a614ec55..4c5e0a61 100644 --- a/webrtc/webrtc.go +++ b/webrtc/webrtc.go @@ -270,7 +270,6 @@ func (w *WebRTC) StopClient() { w.connection.Close() } w.connection = nil - //w.isClosed = true } // IsConnected comment @@ -278,10 +277,6 @@ func (w *WebRTC) IsConnected() bool { return w.isConnected } -func (w *WebRTC) IsClosed() bool { - return w.isClosed -} - // func (w *WebRTC) startStreaming(vp8Track *webrtc.Track, opusTrack *webrtc.Track) { func (w *WebRTC) startStreaming(vp8Track *webrtc.Track, audioTrack *webrtc.DataChannel) { log.Println("Start streaming")