diff --git a/go.mod b/go.mod index b23a3619..d3d550b7 100644 --- a/go.mod +++ b/go.mod @@ -5,7 +5,7 @@ go 1.15 require ( github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751 // indirect github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d // indirect - github.com/aler9/gortsplib v0.0.0-20201115164941-3f45d21f1132 + github.com/aler9/gortsplib v0.0.0-20201115231025-c24d3f677b43 github.com/davecgh/go-spew v1.1.1 // indirect github.com/fsnotify/fsnotify v1.4.9 github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51 diff --git a/go.sum b/go.sum index 2e0a5b4b..408e0fd6 100644 --- a/go.sum +++ b/go.sum @@ -2,8 +2,8 @@ github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751 h1:JYp7IbQjafo github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d h1:UQZhZ2O0vMHr2cI+DC1Mbh0TJxzA3RcLoMsFw+aXw7E= github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho= -github.com/aler9/gortsplib v0.0.0-20201115164941-3f45d21f1132 h1:yPtvZE/tfSPkQn7R2TkBv0I9t318yKDgCwynHzuCK5E= -github.com/aler9/gortsplib v0.0.0-20201115164941-3f45d21f1132/go.mod h1:6yKsTNIrCapRz90WHQtyFV/rKK0TT+QapxUXNqSJi9M= +github.com/aler9/gortsplib v0.0.0-20201115231025-c24d3f677b43 h1:Y+Cl50nLvue3T56uRB1wnq5yLbHpxsjwTjfBZMfA3Qs= +github.com/aler9/gortsplib v0.0.0-20201115231025-c24d3f677b43/go.mod h1:6yKsTNIrCapRz90WHQtyFV/rKK0TT+QapxUXNqSJi9M= github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= diff --git a/internal/client/client.go b/internal/client/client.go index a88f2e98..ffab9724 100644 --- a/internal/client/client.go +++ b/internal/client/client.go @@ -884,25 +884,25 @@ func (c *Client) handleRequest(req *base.Request) error { } func (c *Client) runInitial() bool { - readDone := make(chan error) + readerDone := make(chan error) go func() { for { req, err := c.conn.ReadRequest() if err != nil { - readDone <- err + readerDone <- err return } err = c.handleRequest(req) if err != nil { - readDone <- err + readerDone <- err return } } }() select { - case err := <-readDone: + case err := <-readerDone: switch err { case errStateWaitingDescribe: return c.runWaitingDescribe() @@ -926,7 +926,7 @@ func (c *Client) runInitial() bool { case <-c.terminate: c.conn.Close() - <-readDone + <-readerDone return false } } @@ -1024,25 +1024,25 @@ func (c *Client) runPlay() bool { } func (c *Client) runPlayUDP() bool { - readDone := make(chan error) + readerDone := make(chan error) go func() { for { req, err := c.conn.ReadRequest() if err != nil { - readDone <- err + readerDone <- err return } err = c.handleRequest(req) if err != nil { - readDone <- err + readerDone <- err return } } }() select { - case err := <-readDone: + case err := <-readerDone: if err == errStateInitial { c.state = statePrePlay c.path.OnClientPause(c) @@ -1067,7 +1067,7 @@ func (c *Client) runPlayUDP() bool { c.path = nil c.conn.Close() - <-readDone + <-readerDone return false } } @@ -1076,12 +1076,12 @@ func (c *Client) runPlayTCP() bool { readRequest := make(chan readRequestPair) defer close(readRequest) - readDone := make(chan error) + readerDone := make(chan error) go func() { for { recv, err := c.conn.ReadFrameTCPOrRequest(false) if err != nil { - readDone <- err + readerDone <- err return } @@ -1094,7 +1094,7 @@ func (c *Client) runPlayTCP() bool { readRequest <- readRequestPair{recvt, res} err := <-res if err != nil { - readDone <- err + readerDone <- err return } } @@ -1107,7 +1107,7 @@ func (c *Client) runPlayTCP() bool { case req := <-readRequest: req.res <- c.handleRequest(req.req) - case err := <-readDone: + case err := <-readerDone: if err == errStateInitial { ch := c.tcpFrame go func() { @@ -1165,7 +1165,7 @@ func (c *Client) runPlayTCP() bool { close(c.tcpFrame) c.conn.Close() - <-readDone + <-readerDone return false } } @@ -1242,18 +1242,18 @@ func (c *Client) runRecord() bool { } func (c *Client) runRecordUDP() bool { - readDone := make(chan error) + readerDone := make(chan error) go func() { for { req, err := c.conn.ReadRequest() if err != nil { - readDone <- err + readerDone <- err return } err = c.handleRequest(req) if err != nil { - readDone <- err + readerDone <- err return } } @@ -1267,7 +1267,7 @@ func (c *Client) runRecordUDP() bool { for { select { - case err := <-readDone: + case err := <-readerDone: if err == errStateInitial { for _, track := range c.streamTracks { c.serverUdpRtp.RemovePublisher(c.ip(), track.rtpPort, c) @@ -1312,7 +1312,7 @@ func (c *Client) runRecordUDP() bool { c.log("ERR: no packets received recently (maybe there's a firewall/NAT in between)") c.conn.Close() - <-readDone + <-readerDone c.path.OnClientRemove(c) c.path = nil @@ -1341,7 +1341,7 @@ func (c *Client) runRecordUDP() bool { } c.conn.Close() - <-readDone + <-readerDone c.path.OnClientRemove(c) c.path = nil @@ -1355,19 +1355,19 @@ func (c *Client) runRecordTCP() bool { readRequest := make(chan readRequestPair) defer close(readRequest) - readDone := make(chan error) + readerDone := make(chan error) go func() { for { recv, err := c.conn.ReadFrameTCPOrRequest(true) if err != nil { - readDone <- err + readerDone <- err return } switch recvt := recv.(type) { case *base.InterleavedFrame: if recvt.TrackId >= len(c.streamTracks) { - readDone <- fmt.Errorf("invalid track id '%d'", recvt.TrackId) + readerDone <- fmt.Errorf("invalid track id '%d'", recvt.TrackId) return } @@ -1377,7 +1377,7 @@ func (c *Client) runRecordTCP() bool { case *base.Request: err := c.handleRequest(recvt) if err != nil { - readDone <- err + readerDone <- err return } } @@ -1393,7 +1393,7 @@ func (c *Client) runRecordTCP() bool { case req := <-readRequest: req.res <- c.handleRequest(req.req) - case err := <-readDone: + case err := <-readerDone: if err == errStateInitial { c.state = statePreRecord c.path.OnClientPause(c) @@ -1429,7 +1429,7 @@ func (c *Client) runRecordTCP() bool { }() c.conn.Close() - <-readDone + <-readerDone c.path.OnClientRemove(c) c.path = nil diff --git a/internal/serverudp/server.go b/internal/serverudp/server.go index 3ec30983..cb708dc7 100644 --- a/internal/serverudp/server.go +++ b/internal/serverudp/server.go @@ -111,9 +111,9 @@ func (s *Server) Close() { func (s *Server) run() { defer close(s.done) - writeDone := make(chan struct{}) + writerDone := make(chan struct{}) go func() { - defer close(writeDone) + defer close(writerDone) for w := range s.write { s.pc.SetWriteDeadline(time.Now().Add(s.writeTimeout)) s.pc.WriteTo(w.buf, w.addr) @@ -144,7 +144,7 @@ func (s *Server) run() { } close(s.write) - <-writeDone + <-writerDone } // Port returns the server local port. diff --git a/internal/sourcertmp/source.go b/internal/sourcertmp/source.go index 5f7b7b83..fab89aa6 100644 --- a/internal/sourcertmp/source.go +++ b/internal/sourcertmp/source.go @@ -229,33 +229,33 @@ func (s *Source) runInner() bool { s.parent.Log("rtmp source ready") s.parent.OnSourceSetReady(tracks) - readDone := make(chan error) + readerDone := make(chan error) go func() { for { pkt, err := conn.ReadPacket() if err != nil { - readDone <- err + readerDone <- err return } switch pkt.Type { case av.H264: if h264Sps == nil { - readDone <- fmt.Errorf("rtmp source ERR: received an H264 frame, but track is not setup up") + readerDone <- fmt.Errorf("rtmp source ERR: received an H264 frame, but track is not setup up") return } // decode from AVCC format nalus, typ := h264.SplitNALUs(pkt.Data) if typ != h264.NALU_AVCC { - readDone <- fmt.Errorf("invalid NALU format (%d)", typ) + readerDone <- fmt.Errorf("invalid NALU format (%d)", typ) return } // encode into RTP/H264 format frames, err := h264Encoder.Write(nalus, pkt.Time+pkt.CTime) if err != nil { - readDone <- err + readerDone <- err return } @@ -265,13 +265,13 @@ func (s *Source) runInner() bool { case av.AAC: if aacConfig == nil { - readDone <- fmt.Errorf("rtmp source ERR: received an AAC frame, but track is not setup up") + readerDone <- fmt.Errorf("rtmp source ERR: received an AAC frame, but track is not setup up") return } frames, err := aacEncoder.Write(pkt.Data, pkt.Time+pkt.CTime) if err != nil { - readDone <- err + readerDone <- err return } @@ -280,7 +280,7 @@ func (s *Source) runInner() bool { } default: - readDone <- fmt.Errorf("rtmp source ERR: unexpected packet: %v", pkt.Type) + readerDone <- fmt.Errorf("rtmp source ERR: unexpected packet: %v", pkt.Type) return } } @@ -293,11 +293,11 @@ outer: select { case <-s.terminate: nconn.Close() - <-readDone + <-readerDone ret = false break outer - case err := <-readDone: + case err := <-readerDone: nconn.Close() s.parent.Log("rtmp source ERR: %s", err) ret = true diff --git a/internal/sourcertsp/source.go b/internal/sourcertsp/source.go index 80095324..6de12c0b 100644 --- a/internal/sourcertsp/source.go +++ b/internal/sourcertsp/source.go @@ -139,18 +139,16 @@ func (s *Source) runInner() bool { s.parent.Log("rtsp source ready") s.parent.OnSourceSetReady(tracks) - readDone := make(chan error) - go func() { - for { - trackId, streamType, content, err := conn.ReadFrame() - if err != nil { - readDone <- err - return - } + readerDone := make(chan error) - s.parent.OnFrame(trackId, streamType, content) + conn.OnFrame(func(trackId int, streamType gortsplib.StreamType, content []byte, err error) { + if err != nil { + readerDone <- err + return } - }() + + s.parent.OnFrame(trackId, streamType, content) + }) var ret bool @@ -159,11 +157,11 @@ outer: select { case <-s.terminate: conn.Close() - <-readDone + <-readerDone ret = false break outer - case err := <-readDone: + case err := <-readerDone: conn.Close() s.parent.Log("rtsp source ERR: %s", err) ret = true