Browse Source

remove redundant stream check mechanism in case of tcp publishers

pull/52/head
aler9 6 years ago
parent
commit
efa31937c9
  1. 2
      go.mod
  2. 4
      go.sum
  3. 2
      main.go
  4. 116
      server-client.go

2
go.mod

@ -5,7 +5,7 @@ go 1.12
require ( require (
github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751 // indirect github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751 // indirect
github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d // indirect github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d // indirect
github.com/aler9/gortsplib v0.0.0-20200719202520-de32b1f15ecb github.com/aler9/gortsplib v0.0.0-20200728125613-6edf6a8f9e09
github.com/aler9/sdp/v3 v3.0.0-20200719093237-2c3d108a7436 github.com/aler9/sdp/v3 v3.0.0-20200719093237-2c3d108a7436
github.com/stretchr/testify v1.6.1 github.com/stretchr/testify v1.6.1
gopkg.in/alecthomas/kingpin.v2 v2.2.6 gopkg.in/alecthomas/kingpin.v2 v2.2.6

4
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/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 h1:UQZhZ2O0vMHr2cI+DC1Mbh0TJxzA3RcLoMsFw+aXw7E=
github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho= github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho=
github.com/aler9/gortsplib v0.0.0-20200719202520-de32b1f15ecb h1:R9F835QLbnfLQrOoHZULCrASRC23287Lb6v5LpOt0TY= github.com/aler9/gortsplib v0.0.0-20200728125613-6edf6a8f9e09 h1:Oqs9cVlb/cgeh/jDU/thamzvHESb3cjy04vgGXo4we0=
github.com/aler9/gortsplib v0.0.0-20200719202520-de32b1f15ecb/go.mod h1:kBMvjIdOHRjLdV+oT28JD72JUPpJuwxOc9u72GG8GpY= github.com/aler9/gortsplib v0.0.0-20200728125613-6edf6a8f9e09/go.mod h1:kBMvjIdOHRjLdV+oT28JD72JUPpJuwxOc9u72GG8GpY=
github.com/aler9/sdp/v3 v3.0.0-20200719093237-2c3d108a7436 h1:W0iNErWKvSAyJBNVx+qQoyFrWOFVgS6f/WEME/D3EZc= github.com/aler9/sdp/v3 v3.0.0-20200719093237-2c3d108a7436 h1:W0iNErWKvSAyJBNVx+qQoyFrWOFVgS6f/WEME/D3EZc=
github.com/aler9/sdp/v3 v3.0.0-20200719093237-2c3d108a7436/go.mod h1:OnlEK3QI7YtM+ShZWtGajmOHLZ3bjU80AcIS5e34i1U= github.com/aler9/sdp/v3 v3.0.0-20200719093237-2c3d108a7436/go.mod h1:OnlEK3QI7YtM+ShZWtGajmOHLZ3bjU80AcIS5e34i1U=
github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8= github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8=

2
main.go

@ -538,7 +538,7 @@ func (p *program) forwardFrame(path string, trackId int, streamType gortsplib.St
func main() { func main() {
_, err := newProgram(os.Args[1:], os.Stdin) _, err := newProgram(os.Args[1:], os.Stdin)
if err != nil { if err != nil {
log.Fatal("ERR: ", err) log.Fatal("ERR:", err)
} }
select {} select {}

116
server-client.go

@ -917,43 +917,20 @@ func (c *serverClient) runRecord(path string) {
} }
} }
if c.streamProtocol == gortsplib.StreamProtocolTcp { if c.streamProtocol == gortsplib.StreamProtocolUdp {
frame := &gortsplib.InterleavedFrame{}
readDone := make(chan error) readDone := make(chan error)
go func() { go func() {
for { for {
frame.Content = c.readBuf.swap() req, err := c.conn.ReadRequest()
frame.Content = frame.Content[:cap(frame.Content)]
recv, err := c.conn.ReadFrameOrRequest(frame)
if err != nil { if err != nil {
readDone <- err readDone <- err
break break
} }
switch recvt := recv.(type) { ok := c.handleRequest(req)
case *gortsplib.InterleavedFrame: if !ok {
if frame.TrackId >= len(c.streamTracks) { readDone <- nil
c.log("ERR: invalid track id '%d'", frame.TrackId) break
readDone <- nil
break
}
c.rtcpReceivers[frame.TrackId].OnFrame(frame.StreamType, frame.Content)
c.p.events <- programEventClientFrameTcp{
c.path,
frame.TrackId,
frame.StreamType,
frame.Content,
}
case *gortsplib.Request:
ok := c.handleRequest(recvt)
if !ok {
readDone <- nil
break
}
} }
} }
}() }()
@ -961,14 +938,14 @@ func (c *serverClient) runRecord(path string) {
checkStreamTicker := time.NewTicker(clientCheckStreamInterval) checkStreamTicker := time.NewTicker(clientCheckStreamInterval)
receiverReportTicker := time.NewTicker(clientReceiverReportInterval) receiverReportTicker := time.NewTicker(clientReceiverReportInterval)
outer1: outer2:
for { for {
select { select {
case err := <-readDone: case err := <-readDone:
if err != nil && err != io.EOF { if err != nil && err != io.EOF {
c.log("ERR: %s", err) c.log("ERR: %s", err)
} }
break outer1 break outer2
case <-checkStreamTicker.C: case <-checkStreamTicker.C:
for trackId := range c.streamTracks { for trackId := range c.streamTracks {
@ -976,18 +953,21 @@ func (c *serverClient) runRecord(path string) {
c.log("ERR: stream is dead") c.log("ERR: stream is dead")
c.conn.NetConn().Close() c.conn.NetConn().Close()
<-readDone <-readDone
break outer1 break outer2
} }
} }
case <-receiverReportTicker.C: case <-receiverReportTicker.C:
for trackId := range c.streamTracks { for trackId := range c.streamTracks {
frame := c.rtcpReceivers[trackId].Report() frame := c.rtcpReceivers[trackId].Report()
c.conn.WriteFrame(&gortsplib.InterleavedFrame{ c.p.rtcpl.writeChan <- &udpAddrBufPair{
TrackId: trackId, addr: &net.UDPAddr{
StreamType: gortsplib.StreamTypeRtcp, IP: c.ip(),
Content: frame, Zone: c.zone(),
}) Port: c.streamTracks[trackId].rtcpPort,
},
buf: frame,
}
} }
} }
} }
@ -996,61 +976,69 @@ func (c *serverClient) runRecord(path string) {
receiverReportTicker.Stop() receiverReportTicker.Stop()
} else { } else {
frame := &gortsplib.InterleavedFrame{}
readDone := make(chan error) readDone := make(chan error)
go func() { go func() {
for { for {
req, err := c.conn.ReadRequest() frame.Content = c.readBuf.swap()
frame.Content = frame.Content[:cap(frame.Content)]
recv, err := c.conn.ReadFrameOrRequest(frame)
if err != nil { if err != nil {
readDone <- err readDone <- err
break break
} }
ok := c.handleRequest(req) switch recvt := recv.(type) {
if !ok { case *gortsplib.InterleavedFrame:
readDone <- nil if frame.TrackId >= len(c.streamTracks) {
break c.log("ERR: invalid track id '%d'", frame.TrackId)
readDone <- nil
break
}
c.rtcpReceivers[frame.TrackId].OnFrame(frame.StreamType, frame.Content)
c.p.events <- programEventClientFrameTcp{
c.path,
frame.TrackId,
frame.StreamType,
frame.Content,
}
case *gortsplib.Request:
ok := c.handleRequest(recvt)
if !ok {
readDone <- nil
break
}
} }
} }
}() }()
checkStreamTicker := time.NewTicker(clientCheckStreamInterval)
receiverReportTicker := time.NewTicker(clientReceiverReportInterval) receiverReportTicker := time.NewTicker(clientReceiverReportInterval)
outer2: outer1:
for { for {
select { select {
case err := <-readDone: case err := <-readDone:
if err != nil && err != io.EOF { if err != nil && err != io.EOF {
c.log("ERR: %s", err) c.log("ERR: %s", err)
} }
break outer2 break outer1
case <-checkStreamTicker.C:
for trackId := range c.streamTracks {
if time.Since(c.rtcpReceivers[trackId].LastFrameTime()) >= c.p.conf.ReadTimeout {
c.log("ERR: stream is dead")
c.conn.NetConn().Close()
<-readDone
break outer2
}
}
case <-receiverReportTicker.C: case <-receiverReportTicker.C:
for trackId := range c.streamTracks { for trackId := range c.streamTracks {
frame := c.rtcpReceivers[trackId].Report() frame := c.rtcpReceivers[trackId].Report()
c.p.rtcpl.writeChan <- &udpAddrBufPair{ c.conn.WriteFrame(&gortsplib.InterleavedFrame{
addr: &net.UDPAddr{ TrackId: trackId,
IP: c.ip(), StreamType: gortsplib.StreamTypeRtcp,
Zone: c.zone(), Content: frame,
Port: c.streamTracks[trackId].rtcpPort, })
},
buf: frame,
}
} }
} }
} }
checkStreamTicker.Stop()
receiverReportTicker.Stop() receiverReportTicker.Stop()
} }

Loading…
Cancel
Save