Browse Source

RTMP client: speed up video reading by 1 frame

pull/344/head
aler9 5 years ago
parent
commit
fb0122ba18
  1. 2
      go.mod
  2. 4
      go.sum
  3. 63
      internal/clientrtmp/client.go
  4. 40
      internal/h264/nalutype.go
  5. 26
      internal/sourcertmp/source.go

2
go.mod

@ -5,7 +5,7 @@ go 1.15
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-20210404174403-c2d5ced43bcf github.com/aler9/gortsplib v0.0.0-20210405155604-491b36cdc4b2
github.com/davecgh/go-spew v1.1.1 // indirect github.com/davecgh/go-spew v1.1.1 // indirect
github.com/fsnotify/fsnotify v1.4.9 github.com/fsnotify/fsnotify v1.4.9
github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51 github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51

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-20210404174403-c2d5ced43bcf h1:X8ToQ1H9I7gLp37fxH7mo/U0y+GgF+kwmHpeht4MUHE= github.com/aler9/gortsplib v0.0.0-20210405155604-491b36cdc4b2 h1:AYmBhFE5DTGoQ8XKqmklWO3+CCYYR1X18XL/b/sndmM=
github.com/aler9/gortsplib v0.0.0-20210404174403-c2d5ced43bcf/go.mod h1:zVCg+TQX445hh1pC5QgAuuBvvXZMWLY1XYz626dGFqY= github.com/aler9/gortsplib v0.0.0-20210405155604-491b36cdc4b2/go.mod h1:zVCg+TQX445hh1pC5QgAuuBvvXZMWLY1XYz626dGFqY=
github.com/aler9/rtmp v0.0.0-20210403095203-3be4a5535927 h1:95mXJ5fUCYpBRdSOnLAQAdJHHKxxxJrVCiaqDi965YQ= github.com/aler9/rtmp v0.0.0-20210403095203-3be4a5535927 h1:95mXJ5fUCYpBRdSOnLAQAdJHHKxxxJrVCiaqDi965YQ=
github.com/aler9/rtmp v0.0.0-20210403095203-3be4a5535927/go.mod h1:vzuE21rowz+lT1NGsWbreIvYulgBpCGnQyeTyFblUHc= github.com/aler9/rtmp v0.0.0-20210403095203-3be4a5535927/go.mod h1:vzuE21rowz+lT1NGsWbreIvYulgBpCGnQyeTyFblUHc=
github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8= github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8=

63
internal/clientrtmp/client.go

@ -223,6 +223,7 @@ func (c *Client) runRead() {
var videoTrack *gortsplib.Track var videoTrack *gortsplib.Track
var h264Decoder *rtph264.Decoder var h264Decoder *rtph264.Decoder
var audioTrack *gortsplib.Track var audioTrack *gortsplib.Track
var audioClockRate int
var aacDecoder *rtpaac.Decoder var aacDecoder *rtpaac.Decoder
err = func() error { err = func() error {
@ -241,8 +242,8 @@ func (c *Client) runRead() {
} }
audioTrack = t audioTrack = t
clockRate, _ := audioTrack.ClockRate() audioClockRate, _ = audioTrack.ClockRate()
aacDecoder = rtpaac.NewDecoder(clockRate) aacDecoder = rtpaac.NewDecoder(audioClockRate)
} }
} }
@ -284,9 +285,8 @@ func (c *Client) runRead() {
go func() { go func() {
writerDone <- func() error { writerDone <- func() error {
videoInitialized := false videoInitialized := false
var videoStartDTS time.Time
var videoBuf [][]byte var videoBuf [][]byte
var videoPTS time.Duration var videoStartDTS time.Time
for { for {
data, ok := c.ringBuffer.Pull() data, ok := c.ringBuffer.Pull()
@ -295,10 +295,8 @@ func (c *Client) runRead() {
} }
pair := data.(trackIDPayloadPair) pair := data.(trackIDPayloadPair)
now := time.Now()
if videoTrack != nil && pair.trackID == videoTrack.ID { if videoTrack != nil && pair.trackID == videoTrack.ID {
nts, err := h264Decoder.Decode(pair.buf) nalus, _, err := h264Decoder.Decode(pair.buf)
if err != nil { if err != nil {
if err != rtph264.ErrMorePacketsNeeded { if err != rtph264.ErrMorePacketsNeeded {
c.log(logger.Warn, "unable to decode video track: %v", err) c.log(logger.Warn, "unable to decode video track: %v", err)
@ -306,46 +304,44 @@ func (c *Client) runRead() {
continue continue
} }
for _, nt := range nts { if !videoInitialized {
videoInitialized = true
videoStartDTS = time.Now()
}
for _, nalu := range nalus {
// remove SPS, PPS and AUD, not needed by RTMP // remove SPS, PPS and AUD, not needed by RTMP
typ := h264.NALUType(nt.NALU[0] & 0x1F) typ := h264.NALUType(nalu[0] & 0x1F)
switch typ { switch typ {
case h264.NALUTypeSPS, h264.NALUTypePPS, h264.NALUTypeAccessUnitDelimiter: case h264.NALUTypeSPS, h264.NALUTypePPS, h264.NALUTypeAccessUnitDelimiter:
continue continue
} }
if !videoInitialized { videoBuf = append(videoBuf, nalu)
videoInitialized = true
videoStartDTS = now
videoPTS = nt.Timestamp
} }
// aggregate NALUs by PTS. // RTP marker means that all the NALUs with the same PTS have been received.
// for instance, aggregate a SEI and a IDR. // send them together.
// this delays the stream by one frame, but is required by RTMP. marker := (pair.buf[1] >> 7 & 0x1) > 0
if nt.Timestamp != videoPTS { if marker {
c.conn.NetConn().SetWriteDeadline(time.Now().Add(c.writeTimeout)) c.conn.NetConn().SetWriteDeadline(time.Now().Add(c.writeTimeout))
err := c.conn.WriteH264(videoBuf, now.Sub(videoStartDTS)) err := c.conn.WriteH264(videoBuf, time.Since(videoStartDTS))
if err != nil { if err != nil {
return err return err
} }
videoBuf = nil videoBuf = nil
} }
videoPTS = nt.Timestamp
videoBuf = append(videoBuf, nt.NALU)
}
} else if audioTrack != nil && pair.trackID == audioTrack.ID { } else if audioTrack != nil && pair.trackID == audioTrack.ID {
ats, err := aacDecoder.Decode(pair.buf) aus, pts, err := aacDecoder.Decode(pair.buf)
if err != nil { if err != nil {
c.log(logger.Warn, "unable to decode audio track: %v", err) c.log(logger.Warn, "unable to decode audio track: %v", err)
continue continue
} }
for _, at := range ats { for i, au := range aus {
c.conn.NetConn().SetWriteDeadline(time.Now().Add(c.writeTimeout)) c.conn.NetConn().SetWriteDeadline(time.Now().Add(c.writeTimeout))
err := c.conn.WriteAAC(at.AU, at.Timestamp) err := c.conn.WriteAAC(au, pts+time.Duration(i)*1000*time.Second/time.Duration(audioClockRate))
if err != nil { if err != nil {
return err return err
} }
@ -520,8 +516,7 @@ func (c *Client) runPublish() {
return err return err
} }
ts := pkt.Time + pkt.CTime var outNALUs [][]byte
var nts []*rtph264.NALUAndTimestamp
for _, nalu := range nalus { for _, nalu := range nalus {
// remove SPS, PPS and AUD, not needed by RTSP / RTMP // remove SPS, PPS and AUD, not needed by RTSP / RTMP
typ := h264.NALUType(nalu[0] & 0x1F) typ := h264.NALUType(nalu[0] & 0x1F)
@ -530,13 +525,10 @@ func (c *Client) runPublish() {
continue continue
} }
nts = append(nts, &rtph264.NALUAndTimestamp{ outNALUs = append(outNALUs, nalu)
Timestamp: ts,
NALU: nalu,
})
} }
frames, err := h264Encoder.Encode(nts) frames, err := h264Encoder.Encode(outNALUs, pkt.Time+pkt.CTime)
if err != nil { if err != nil {
return fmt.Errorf("ERR while encoding H264: %v", err) return fmt.Errorf("ERR while encoding H264: %v", err)
} }
@ -550,12 +542,7 @@ func (c *Client) runPublish() {
return fmt.Errorf("ERR: received an AAC frame, but track is not set up") return fmt.Errorf("ERR: received an AAC frame, but track is not set up")
} }
frames, err := aacEncoder.Encode([]*rtpaac.AUAndTimestamp{ frames, err := aacEncoder.Encode([][]byte{pkt.Data}, pkt.Time+pkt.CTime)
{
Timestamp: pkt.Time + pkt.CTime,
AU: pkt.Data,
},
})
if err != nil { if err != nil {
return fmt.Errorf("ERR while encoding AAC: %v", err) return fmt.Errorf("ERR while encoding AAC: %v", err)
} }

40
internal/h264/nalutype.go

@ -14,7 +14,7 @@ const (
NALUTypeDataPartitionB NALUType = 3 NALUTypeDataPartitionB NALUType = 3
NALUTypeDataPartitionC NALUType = 4 NALUTypeDataPartitionC NALUType = 4
NALUTypeIDR NALUType = 5 NALUTypeIDR NALUType = 5
NALUTypeSei NALUType = 6 NALUTypeSEI NALUType = 6
NALUTypeSPS NALUType = 7 NALUTypeSPS NALUType = 7
NALUTypePPS NALUType = 8 NALUTypePPS NALUType = 8
NALUTypeAccessUnitDelimiter NALUType = 9 NALUTypeAccessUnitDelimiter NALUType = 9
@ -32,12 +32,12 @@ const (
NALUTypeSliceExtensionDepth NALUType = 21 NALUTypeSliceExtensionDepth NALUType = 21
NALUTypeReserved22 NALUType = 22 NALUTypeReserved22 NALUType = 22
NALUTypeReserved23 NALUType = 23 NALUTypeReserved23 NALUType = 23
NALUTypeStapA NALUType = 24 NALUTypeSTAPA NALUType = 24
NALUTypeStapB NALUType = 25 NALUTypeSTAPB NALUType = 25
NALUTypeMtap16 NALUType = 26 NALUTypeMTAP16 NALUType = 26
NALUTypeMtap24 NALUType = 27 NALUTypeMTAP24 NALUType = 27
NALUTypeFuA NALUType = 28 NALUTypeFUA NALUType = 28
NALUTypeFuB NALUType = 29 NALUTypeFUB NALUType = 29
) )
// String implements fmt.Stringer. // String implements fmt.Stringer.
@ -53,7 +53,7 @@ func (nt NALUType) String() string {
return "DataPartitionC" return "DataPartitionC"
case NALUTypeIDR: case NALUTypeIDR:
return "IDR" return "IDR"
case NALUTypeSei: case NALUTypeSEI:
return "Sei" return "Sei"
case NALUTypeSPS: case NALUTypeSPS:
return "SPS" return "SPS"
@ -89,18 +89,18 @@ func (nt NALUType) String() string {
return "Reserved22" return "Reserved22"
case NALUTypeReserved23: case NALUTypeReserved23:
return "Reserved23" return "Reserved23"
case NALUTypeStapA: case NALUTypeSTAPA:
return "StapA" return "STAPA"
case NALUTypeStapB: case NALUTypeSTAPB:
return "StapB" return "STAPB"
case NALUTypeMtap16: case NALUTypeMTAP16:
return "Mtap16" return "MTAP16"
case NALUTypeMtap24: case NALUTypeMTAP24:
return "Mtap24" return "MTAP24"
case NALUTypeFuA: case NALUTypeFUA:
return "FuA" return "FUA"
case NALUTypeFuB: case NALUTypeFUB:
return "FuB" return "FUB"
} }
return fmt.Sprintf("unknown (%d)", nt) return fmt.Sprintf("unknown (%d)", nt)
} }

26
internal/sourcertmp/source.go

@ -214,16 +214,19 @@ func (s *Source) runInner() bool {
return err return err
} }
ts := pkt.Time + pkt.CTime var outNALUs [][]byte
var nts []*rtph264.NALUAndTimestamp for _, nalu := range nalus {
for _, nt := range nalus { // remove SPS, PPS and AUD, not needed by RTSP / RTMP
nts = append(nts, &rtph264.NALUAndTimestamp{ typ := h264.NALUType(nalu[0] & 0x1F)
Timestamp: ts, switch typ {
NALU: nt, case h264.NALUTypeSPS, h264.NALUTypePPS, h264.NALUTypeAccessUnitDelimiter:
}) continue
}
outNALUs = append(outNALUs, nalu)
} }
frames, err := h264Encoder.Encode(nts) frames, err := h264Encoder.Encode(outNALUs, pkt.Time+pkt.CTime)
if err != nil { if err != nil {
return fmt.Errorf("ERR while encoding H264: %v", err) return fmt.Errorf("ERR while encoding H264: %v", err)
} }
@ -237,12 +240,7 @@ func (s *Source) runInner() bool {
return fmt.Errorf("ERR: received an AAC frame, but track is not set up") return fmt.Errorf("ERR: received an AAC frame, but track is not set up")
} }
frames, err := aacEncoder.Encode([]*rtpaac.AUAndTimestamp{ frames, err := aacEncoder.Encode([][]byte{pkt.Data}, pkt.Time+pkt.CTime)
{
Timestamp: pkt.Time + pkt.CTime,
AU: pkt.Data,
},
})
if err != nil { if err != nil {
return fmt.Errorf("ERR while encoding AAC: %v", err) return fmt.Errorf("ERR while encoding AAC: %v", err)
} }

Loading…
Cancel
Save