From c75efdc559fc3b6d0ac2f46d9d65cacaeb3e06ee Mon Sep 17 00:00:00 2001 From: giongto35 Date: Tue, 14 May 2019 03:03:51 +0800 Subject: [PATCH] Update close encoder channel --- handler/handlers.go | 6 +++ handler/room/room.go | 2 +- vpx-encoder/encoder.go | 110 ++++++++++++++++++++++++----------------- webrtc/webrtc.go | 5 +- 4 files changed, 75 insertions(+), 48 deletions(-) diff --git a/handler/handlers.go b/handler/handlers.go index c2a536f4..fd54d663 100644 --- a/handler/handlers.go +++ b/handler/handlers.go @@ -74,6 +74,12 @@ func (h *Handler) GetWeb(w http.ResponseWriter, r *http.Request) { // WS handles normal traffic (from browser to host) func (h *Handler) WS(w http.ResponseWriter, r *http.Request) { + defer func() { + if r := recover(); r != nil { + log.Println("Warn: Something wrong. Recovered in f", r) + } + }() + c, err := upgrader.Upgrade(w, r, nil) if err != nil { log.Print("[!] WS upgrade:", err) diff --git a/handler/room/room.go b/handler/room/room.go index deddca6f..a5b6f5f0 100644 --- a/handler/room/room.go +++ b/handler/room/room.go @@ -110,7 +110,7 @@ func (r *Room) AddConnectionToRoom(peerconnection *webrtc.WebRTC, playerIndex in func (r *Room) startWebRTCSession(peerconnection *webrtc.WebRTC, playerIndex int) { defer func() { if r := recover(); r != nil { - fmt.Println("Recovered when sent to close inputChannel") + log.Println("Warn: Recovered when sent to close inputChannel") } }() diff --git a/vpx-encoder/encoder.go b/vpx-encoder/encoder.go index a76300a0..4cc07998 100644 --- a/vpx-encoder/encoder.go +++ b/vpx-encoder/encoder.go @@ -48,6 +48,9 @@ func NewVpxEncoder(w, h, fps, bitrate, keyframe int) (*VpxEncoder, error) { v := &VpxEncoder{ Output: make(chan []byte, 5*chanSize), Input: make(chan []byte, chanSize), + + IsRunning: true, + Done: false, // C width: C.uint(w), height: C.uint(h), @@ -66,9 +69,11 @@ func NewVpxEncoder(w, h, fps, bitrate, keyframe int) (*VpxEncoder, error) { // VpxEncoder yuvI420 image to vp8 video type VpxEncoder struct { - started bool - Output chan []byte // frame - Input chan []byte // yuvI420 + Output chan []byte // frame + Input chan []byte // yuvI420 + + IsRunning bool + Done bool // C width C.uint height C.uint @@ -109,61 +114,74 @@ func (v *VpxEncoder) init() error { if C.call_vpx_codec_enc_init(&v.vpxCodexCtx, encoder, &cfg) != 0 { return fmt.Errorf("Failed to initialize encoder") } - v.started = true + v.IsRunning = true go v.startLooping() return nil } func (v *VpxEncoder) startLooping() { - go func() { - for { - beginEncoding := time.Now() - - yuv, ok := <-v.Input - if !ok { - return - } - // Add Image - v.vpxCodexIter = nil - C.vpx_img_read(&v.vpxImage, unsafe.Pointer(&yuv[0])) - - var flags C.int - if v.keyFrameInterval > 0 && v.frameCount%v.keyFrameInterval == 0 { - flags |= C.VPX_EFLAG_FORCE_KF - } - if C.vpx_codec_encode(&v.vpxCodexCtx, &v.vpxImage, C.vpx_codec_pts_t(v.frameCount), 1, C.vpx_enc_frame_flags_t(flags), C.VPX_DL_REALTIME) != 0 { - fmt.Println("Failed to encode frame") - } - v.frameCount++ - - // Get Frame - for { - goBytes := C.get_frame_buffer(&v.vpxCodexCtx, &v.vpxCodexIter) - if goBytes.bs == nil { - break - } - bs := C.GoBytes(goBytes.bs, goBytes.size) - // if buffer is full skip frame - if len(v.Output) >= cap(v.Output) { - continue - } - v.Output <- bs - } - - if *config.IsMonitor { - log.Println("Encoding time: ", time.Now().Sub(beginEncoding)) - } + defer func() { + if r := recover(); r != nil { + log.Println("Warn: Recovered panic in encoding ", r) } }() + + for { + beginEncoding := time.Now() + + yuv, ok := <-v.Input + if !ok || v.Done == true { + // The first time we see IsRunning set to false, we release and return + v.Release() + return + } + + // Add Image + v.vpxCodexIter = nil + C.vpx_img_read(&v.vpxImage, unsafe.Pointer(&yuv[0])) + + var flags C.int + if v.keyFrameInterval > 0 && v.frameCount%v.keyFrameInterval == 0 { + flags |= C.VPX_EFLAG_FORCE_KF + } + if C.vpx_codec_encode(&v.vpxCodexCtx, &v.vpxImage, C.vpx_codec_pts_t(v.frameCount), 1, C.vpx_enc_frame_flags_t(flags), C.VPX_DL_REALTIME) != 0 { + fmt.Println("Failed to encode frame") + } + v.frameCount++ + + // Get Frame + for { + goBytes := C.get_frame_buffer(&v.vpxCodexCtx, &v.vpxCodexIter) + if goBytes.bs == nil { + break + } + bs := C.GoBytes(goBytes.bs, goBytes.size) + // if buffer is full skip frame + if len(v.Output) >= cap(v.Output) { + continue + } + v.Output <- bs + } + + if *config.IsMonitor { + log.Println("Encoding time: ", time.Now().Sub(beginEncoding)) + } + } } // Release release memory and stop loop func (v *VpxEncoder) Release() { - if v.started { - close(v.Input) - close(v.Output) + if v.IsRunning { + log.Println("Releasing encoder") + log.Println("Close output", v.Output) C.vpx_img_free(&v.vpxImage) C.vpx_codec_destroy(&v.vpxCodexCtx) + // TODO: Check here, after close it will signal + if v.Input != nil { + close(v.Input) + } + close(v.Output) } - v.started = false + v.IsRunning = false + // TODO: Can we merge IsRunning and Done together } diff --git a/webrtc/webrtc.go b/webrtc/webrtc.go index d29bdab1..f1b0b1cd 100644 --- a/webrtc/webrtc.go +++ b/webrtc/webrtc.go @@ -232,13 +232,16 @@ func (w *WebRTC) StopClient() { log.Println("===StopClient===") w.isConnected = false if w.encoder != nil { - w.encoder.Release() + // NOTE: We signal using bool value instead of channel for better performance + w.encoder.Done = true } if w.connection != nil { w.connection.Close() } w.connection = nil close(w.InputChannel) + // webrtc is producer, so we close + close(w.encoder.Input) // NOTE: ImageChannel is waiting for input. Close in writer is not correct for this close(w.ImageChannel) close(w.AudioChannel)