mediamtx/internal/core/hls_source.go

176 lines
3.1 KiB
Go
Raw Normal View History

2021-09-05 15:35:04 +00:00
package core
import (
"context"
"sync"
"time"
"github.com/aler9/gortsplib"
"github.com/pion/rtp"
2021-09-05 15:35:04 +00:00
"github.com/aler9/rtsp-simple-server/internal/hls"
"github.com/aler9/rtsp-simple-server/internal/logger"
"github.com/aler9/rtsp-simple-server/internal/rtcpsenderset"
)
const (
hlsSourceRetryPause = 5 * time.Second
)
type hlsSourceParent interface {
2021-10-27 19:01:00 +00:00
log(logger.Level, string, ...interface{})
onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes
2022-01-14 22:42:41 +00:00
onSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq)
2021-09-05 15:35:04 +00:00
}
type hlsSource struct {
ur string
fingerprint string
wg *sync.WaitGroup
parent hlsSourceParent
2021-09-05 15:35:04 +00:00
ctx context.Context
ctxCancel func()
}
func newHLSSource(
parentCtx context.Context,
ur string,
fingerprint string,
2021-09-05 15:35:04 +00:00
wg *sync.WaitGroup,
parent hlsSourceParent) *hlsSource {
ctx, ctxCancel := context.WithCancel(parentCtx)
s := &hlsSource{
ur: ur,
fingerprint: fingerprint,
wg: wg,
parent: parent,
ctx: ctx,
ctxCancel: ctxCancel,
2021-09-05 15:35:04 +00:00
}
s.Log(logger.Info, "started")
s.wg.Add(1)
go s.run()
return s
}
2021-10-27 19:01:00 +00:00
func (s *hlsSource) close() {
2021-09-05 15:35:04 +00:00
s.Log(logger.Info, "stopped")
s.ctxCancel()
}
func (s *hlsSource) Log(level logger.Level, format string, args ...interface{}) {
2021-10-27 19:01:00 +00:00
s.parent.log(level, "[hls source] "+format, args...)
2021-09-05 15:35:04 +00:00
}
func (s *hlsSource) run() {
defer s.wg.Done()
outer:
for {
ok := s.runInner()
if !ok {
break outer
}
select {
case <-time.After(hlsSourceRetryPause):
case <-s.ctx.Done():
break outer
}
}
s.ctxCancel()
}
func (s *hlsSource) runInner() bool {
var stream *stream
var rtcpSenders *rtcpsenderset.RTCPSenderSet
var videoTrackID int
var audioTrackID int
defer func() {
if stream != nil {
2022-01-14 22:42:41 +00:00
s.parent.onSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{source: s})
2021-09-05 15:35:04 +00:00
rtcpSenders.Close()
}
}()
2022-01-30 16:36:42 +00:00
onTracks := func(videoTrack gortsplib.Track, audioTrack gortsplib.Track) error {
2021-09-05 15:35:04 +00:00
var tracks gortsplib.Tracks
if videoTrack != nil {
videoTrackID = len(tracks)
tracks = append(tracks, videoTrack)
}
if audioTrack != nil {
audioTrackID = len(tracks)
tracks = append(tracks, audioTrack)
}
2021-10-27 19:01:00 +00:00
res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{
2022-01-14 22:42:41 +00:00
source: s,
tracks: tracks,
2021-09-05 15:35:04 +00:00
})
2022-01-14 22:42:41 +00:00
if res.err != nil {
return res.err
2021-09-05 15:35:04 +00:00
}
s.Log(logger.Info, "ready")
2022-01-14 22:42:41 +00:00
stream = res.stream
2021-11-12 21:29:56 +00:00
rtcpSenders = rtcpsenderset.New(tracks, stream.onPacketRTCP)
2021-09-05 15:35:04 +00:00
return nil
}
onPacket := func(isVideo bool, pkt *rtp.Packet) {
2021-09-05 15:35:04 +00:00
var trackID int
if isVideo {
trackID = videoTrackID
} else {
trackID = audioTrackID
}
if stream != nil {
rtcpSenders.OnPacketRTP(trackID, pkt)
stream.onPacketRTP(trackID, pkt)
2021-09-05 15:35:04 +00:00
}
}
2021-10-28 16:57:04 +00:00
c, err := hls.NewClient(
2021-09-05 15:35:04 +00:00
s.ur,
s.fingerprint,
2021-09-05 15:35:04 +00:00
onTracks,
2021-11-12 21:29:56 +00:00
onPacket,
2021-09-05 15:35:04 +00:00
s,
)
2021-10-28 16:57:04 +00:00
if err != nil {
s.Log(logger.Info, "ERR: %v", err)
return true
}
2021-09-05 15:35:04 +00:00
select {
case err := <-c.Wait():
s.Log(logger.Info, "ERR: %v", err)
return true
case <-s.ctx.Done():
c.Close()
<-c.Wait()
return false
}
}
2021-10-27 19:01:00 +00:00
// onSourceAPIDescribe implements source.
func (*hlsSource) onSourceAPIDescribe() interface{} {
2021-09-05 15:35:04 +00:00
return struct {
Type string `json:"type"`
}{"hlsSource"}
}