Browse Source

speed up RTMP sources

pull/235/head
aler9 6 years ago
parent
commit
49ac52ff67
  1. 296
      internal/sourcertmp/source.go
  2. 12
      main_test.go
  3. 2
      testimages/ffmpeg/Dockerfile
  4. BIN
      testimages/ffmpeg/emptyvideoaudio.ts

296
internal/sourcertmp/source.go

@ -13,6 +13,7 @@ import (
"github.com/aler9/gortsplib/pkg/rtph264" "github.com/aler9/gortsplib/pkg/rtph264"
"github.com/notedit/rtmp/av" "github.com/notedit/rtmp/av"
"github.com/notedit/rtmp/codec/h264" "github.com/notedit/rtmp/codec/h264"
"github.com/notedit/rtmp/format/flv/flvio"
"github.com/notedit/rtmp/format/rtmp" "github.com/notedit/rtmp/format/rtmp"
"github.com/aler9/rtsp-simple-server/internal/logger" "github.com/aler9/rtsp-simple-server/internal/logger"
@ -20,8 +21,10 @@ import (
) )
const ( const (
retryPause = 5 * time.Second retryPause = 5 * time.Second
analyzeTimeout = 8 * time.Second
codecH264 = 7
codecAAC = 10
) )
// Parent is implemeneted by path.Path. // Parent is implemeneted by path.Path.
@ -107,6 +110,33 @@ func (s *Source) run() {
} }
} }
func readMetadata(conn *rtmp.Conn) (flvio.AMFMap, error) {
pkt, err := conn.ReadPacket()
if err != nil {
return nil, err
}
if pkt.Type != av.Metadata {
return nil, fmt.Errorf("first packet must be metadata")
}
arr, err := flvio.ParseAMFVals(pkt.Data, false)
if err != nil {
return nil, err
}
if len(arr) != 1 {
return nil, fmt.Errorf("invalid metadata")
}
ma, ok := arr[0].(flvio.AMFMap)
if !ok {
return nil, fmt.Errorf("invalid metadata")
}
return ma, nil
}
func (s *Source) runInner() bool { func (s *Source) runInner() bool {
s.log(logger.Info, "connecting") s.log(logger.Info, "connecting")
@ -130,115 +160,130 @@ func (s *Source) runInner() bool {
return true return true
} }
// gather video and audio features var tracks gortsplib.Tracks
var h264Sps []byte
var h264Pps []byte
var aacConfig []byte
confDone := make(chan struct{})
confClose := uint32(0)
go func() {
defer close(confDone)
for { var videoTrack *gortsplib.Track
var pkt av.Packet var videoRTCPSender *rtcpsender.RTCPSender
pkt, err = conn.ReadPacket() var h264Encoder *rtph264.Encoder
if err != nil {
return
}
if atomic.LoadUint32(&confClose) > 0 { var audioTrack *gortsplib.Track
return var audioRTCPSender *rtcpsender.RTCPSender
var aacEncoder *rtpaac.Encoder
confDone := make(chan error)
go func() {
confDone <- func() error {
md, err := readMetadata(conn)
if err != nil {
return err
} }
switch pkt.Type { hasVideo := false
case av.H264DecoderConfig: if v, ok := md.GetFloat64("videocodecid"); ok {
codec, err := h264.FromDecoderConfig(pkt.Data) switch v {
if err != nil { case codecH264:
panic(err) hasVideo = true
case 0:
default:
return fmt.Errorf("unsupported video codec %v", v)
} }
h264Sps, h264Pps = codec.SPS[0], codec.PPS[0] }
if aacConfig != nil { hasAudio := false
return if v, ok := md.GetFloat64("audiocodecid"); ok {
switch v {
case codecAAC:
hasAudio = true
case 0:
default:
return fmt.Errorf("unsupported audio codec %v", v)
} }
}
case av.AACDecoderConfig: if !hasVideo && !hasAudio {
aacConfig = pkt.Data return fmt.Errorf("stream has no tracks")
}
if h264Sps != nil { for {
return var pkt av.Packet
pkt, err = conn.ReadPacket()
if err != nil {
return err
} }
}
}
}()
timer := time.NewTimer(analyzeTimeout) switch pkt.Type {
defer timer.Stop() case av.H264DecoderConfig:
if !hasVideo {
return fmt.Errorf("unexpected video packet")
}
if videoTrack != nil {
return fmt.Errorf("video track setupped twice")
}
select { codec, err := h264.FromDecoderConfig(pkt.Data)
case <-confDone: if err != nil {
case <-timer.C: return err
atomic.StoreUint32(&confClose, 1) }
<-confDone
}
if err != nil { videoTrack, err = gortsplib.NewTrackH264(96, codec.SPS[0], codec.PPS[0])
s.log(logger.Info, "ERR: %s", err) if err != nil {
return true return err
} }
var tracks gortsplib.Tracks clockRate, _ := videoTrack.ClockRate()
videoRTCPSender = rtcpsender.New(clockRate)
var videoTrack *gortsplib.Track h264Encoder, err = rtph264.NewEncoder(96)
var videoRTCPSender *rtcpsender.RTCPSender if err != nil {
var h264Encoder *rtph264.Encoder return err
}
var audioTrack *gortsplib.Track tracks = append(tracks, videoTrack)
var audioRTCPSender *rtcpsender.RTCPSender
var aacEncoder *rtpaac.Encoder
if h264Sps != nil { case av.AACDecoderConfig:
videoTrack, err = gortsplib.NewTrackH264(96, h264Sps, h264Pps) if !hasAudio {
if err != nil { return fmt.Errorf("unexpected audio packet")
s.log(logger.Info, "ERR: %s", err) }
return true if audioTrack != nil {
} return fmt.Errorf("audio track setupped twice")
}
clockRate, _ := videoTrack.ClockRate() audioTrack, err = gortsplib.NewTrackAAC(96, pkt.Data)
videoRTCPSender = rtcpsender.New(clockRate) if err != nil {
return err
}
h264Encoder, err = rtph264.NewEncoder(96) clockRate, _ := audioTrack.ClockRate()
if err != nil { audioRTCPSender = rtcpsender.New(clockRate)
s.log(logger.Info, "ERR: %s", err)
return true
}
tracks = append(tracks, videoTrack) aacEncoder, err = rtpaac.NewEncoder(96, clockRate)
} if err != nil {
return err
}
if aacConfig != nil { tracks = append(tracks, audioTrack)
audioTrack, err = gortsplib.NewTrackAAC(96, aacConfig) }
if err != nil {
s.log(logger.Info, "ERR: %s", err)
return true
}
clockRate, _ := audioTrack.ClockRate() if (!hasVideo || videoTrack != nil) &&
audioRTCPSender = rtcpsender.New(clockRate) (!hasAudio || audioTrack != nil) {
return nil
}
}
}()
}()
aacEncoder, err = rtpaac.NewEncoder(96, clockRate) select {
case err := <-confDone:
if err != nil { if err != nil {
s.log(logger.Info, "ERR: %s", err) s.log(logger.Info, "ERR: %s", err)
return true return true
} }
tracks = append(tracks, audioTrack) case <-s.terminate:
} nconn.Close()
<-confDone
if len(tracks) == 0 { return false
s.log(logger.Info, "ERR: no tracks found")
return true
} }
for i, t := range tracks { for i, t := range tracks {
@ -284,61 +329,56 @@ func (s *Source) runInner() bool {
readerDone := make(chan error) readerDone := make(chan error)
go func() { go func() {
for { readerDone <- func() error {
pkt, err := conn.ReadPacket() for {
if err != nil { pkt, err := conn.ReadPacket()
readerDone <- err if err != nil {
return return err
}
switch pkt.Type {
case av.H264:
if h264Sps == nil {
readerDone <- fmt.Errorf("rtmp source ERR: received an H264 frame, but track is not setup up")
return
} }
// decode from AVCC format switch pkt.Type {
nalus, typ := h264.SplitNALUs(pkt.Data) case av.H264:
if typ != h264.NALU_AVCC { if videoTrack == nil {
readerDone <- fmt.Errorf("invalid NALU format (%d)", typ) return fmt.Errorf("rtmp source ERR: received an H264 frame, but track is not setup up")
return }
}
// encode into RTP/H264 format // decode from AVCC format
frames, err := h264Encoder.Write(pkt.Time+pkt.CTime, nalus) nalus, typ := h264.SplitNALUs(pkt.Data)
if err != nil { if typ != h264.NALU_AVCC {
readerDone <- err return fmt.Errorf("invalid NALU format (%d)", typ)
return }
}
for _, f := range frames { // encode into RTP/H264 format
videoRTCPSender.ProcessFrame(time.Now(), gortsplib.StreamTypeRTP, f) frames, err := h264Encoder.Write(pkt.Time+pkt.CTime, nalus)
s.parent.OnFrame(videoTrack.ID, gortsplib.StreamTypeRTP, f) if err != nil {
} return err
}
case av.AAC: for _, f := range frames {
if aacConfig == nil { videoRTCPSender.ProcessFrame(time.Now(), gortsplib.StreamTypeRTP, f)
readerDone <- fmt.Errorf("rtmp source ERR: received an AAC frame, but track is not setup up") s.parent.OnFrame(videoTrack.ID, gortsplib.StreamTypeRTP, f)
return }
}
frames, err := aacEncoder.Write(pkt.Time+pkt.CTime, pkt.Data) case av.AAC:
if err != nil { if audioTrack == nil {
readerDone <- err return fmt.Errorf("rtmp source ERR: received an AAC frame, but track is not setup up")
return }
}
for _, f := range frames { frames, err := aacEncoder.Write(pkt.Time+pkt.CTime, pkt.Data)
audioRTCPSender.ProcessFrame(time.Now(), gortsplib.StreamTypeRTP, f) if err != nil {
s.parent.OnFrame(audioTrack.ID, gortsplib.StreamTypeRTP, f) return err
} }
default: for _, f := range frames {
readerDone <- fmt.Errorf("rtmp source ERR: unexpected packet: %v", pkt.Type) audioRTCPSender.ProcessFrame(time.Now(), gortsplib.StreamTypeRTP, f)
return s.parent.OnFrame(audioTrack.ID, gortsplib.StreamTypeRTP, f)
}
default:
return fmt.Errorf("rtmp source ERR: unexpected packet: %v", pkt.Type)
}
} }
} }()
}() }()
for { for {

12
main_test.go

@ -651,7 +651,8 @@ func TestSource(t *testing.T) {
"rtsp_udp", "rtsp_udp",
"rtsp_tcp", "rtsp_tcp",
"rtsps", "rtsps",
"rtmp", "rtmp_videoaudio",
"rtmp_video",
} { } {
t.Run(source, func(t *testing.T) { t.Run(source, func(t *testing.T) {
switch source { switch source {
@ -730,15 +731,20 @@ func TestSource(t *testing.T) {
require.Equal(t, true, ok) require.Equal(t, true, ok)
defer p2.close() defer p2.close()
case "rtmp": case "rtmp_videoaudio", "rtmp_video":
cnt1, err := newContainer("nginx-rtmp", "rtmpserver", []string{}) cnt1, err := newContainer("nginx-rtmp", "rtmpserver", []string{})
require.NoError(t, err) require.NoError(t, err)
defer cnt1.close() defer cnt1.close()
input := "emptyvideoaudio.ts"
if source == "rtmp_video" {
input = "emptyvideo.ts"
}
cnt2, err := newContainer("ffmpeg", "source", []string{ cnt2, err := newContainer("ffmpeg", "source", []string{
"-re", "-re",
"-stream_loop", "-1", "-stream_loop", "-1",
"-i", "emptyvideo.ts", "-i", input,
"-c", "copy", "-c", "copy",
"-f", "flv", "-f", "flv",
"rtmp://" + cnt1.ip() + "/stream/test", "rtmp://" + cnt1.ip() + "/stream/test",

2
testimages/ffmpeg/Dockerfile

@ -3,7 +3,7 @@ FROM amd64/alpine:3.12
RUN apk add --no-cache \ RUN apk add --no-cache \
ffmpeg ffmpeg
COPY emptyvideo.ts / COPY *.ts /
COPY start.sh / COPY start.sh /
RUN chmod +x /start.sh RUN chmod +x /start.sh

BIN
testimages/ffmpeg/emptyvideoaudio.ts

Binary file not shown.
Loading…
Cancel
Save