Browse Source

hls: change Muxer letter

pull/707/head
aler9 5 years ago
parent
commit
6564c3511b
  1. 126
      internal/core/hls_muxer.go

126
internal/core/hls_muxer.go

@ -155,7 +155,7 @@ func newHLSMuxer(
parent hlsMuxerParent) *hlsMuxer { parent hlsMuxerParent) *hlsMuxer {
ctx, ctxCancel := context.WithCancel(parentCtx) ctx, ctxCancel := context.WithCancel(parentCtx)
r := &hlsMuxer{ m := &hlsMuxer{
hlsAlwaysRemux: hlsAlwaysRemux, hlsAlwaysRemux: hlsAlwaysRemux,
hlsSegmentCount: hlsSegmentCount, hlsSegmentCount: hlsSegmentCount,
hlsSegmentDuration: hlsSegmentDuration, hlsSegmentDuration: hlsSegmentDuration,
@ -173,35 +173,35 @@ func newHLSMuxer(
request: make(chan hlsMuxerRequest), request: make(chan hlsMuxerRequest),
} }
r.log(logger.Info, "opened") m.log(logger.Info, "opened")
r.wg.Add(1) m.wg.Add(1)
go r.run() go m.run()
return r return m
} }
func (r *hlsMuxer) close() { func (m *hlsMuxer) close() {
r.ctxCancel() m.ctxCancel()
} }
func (r *hlsMuxer) log(level logger.Level, format string, args ...interface{}) { func (m *hlsMuxer) log(level logger.Level, format string, args ...interface{}) {
r.parent.log(level, "[muxer %s] "+format, append([]interface{}{r.pathName}, args...)...) m.parent.log(level, "[muxer %s] "+format, append([]interface{}{m.pathName}, args...)...)
} }
// PathName returns the path name. // PathName returns the path name.
func (r *hlsMuxer) PathName() string { func (m *hlsMuxer) PathName() string {
return r.pathName return m.pathName
} }
func (r *hlsMuxer) run() { func (m *hlsMuxer) run() {
defer r.wg.Done() defer m.wg.Done()
innerCtx, innerCtxCancel := context.WithCancel(context.Background()) innerCtx, innerCtxCancel := context.WithCancel(context.Background())
innerReady := make(chan struct{}) innerReady := make(chan struct{})
innerErr := make(chan error) innerErr := make(chan error)
go func() { go func() {
innerErr <- r.runInner(innerCtx, innerReady) innerErr <- m.runInner(innerCtx, innerReady)
}() }()
isReady := false isReady := false
@ -209,24 +209,24 @@ func (r *hlsMuxer) run() {
err := func() error { err := func() error {
for { for {
select { select {
case <-r.ctx.Done(): case <-m.ctx.Done():
innerCtxCancel() innerCtxCancel()
<-innerErr <-innerErr
return errors.New("terminated") return errors.New("terminated")
case req := <-r.request: case req := <-m.request:
if isReady { if isReady {
req.Res <- r.handleRequest(req) req.Res <- m.handleRequest(req)
} else { } else {
r.requests = append(r.requests, req) m.requests = append(m.requests, req)
} }
case <-innerReady: case <-innerReady:
isReady = true isReady = true
for _, req := range r.requests { for _, req := range m.requests {
req.Res <- r.handleRequest(req) req.Res <- m.handleRequest(req)
} }
r.requests = nil m.requests = nil
case err := <-innerErr: case err := <-innerErr:
innerCtxCancel() innerCtxCancel()
@ -235,21 +235,21 @@ func (r *hlsMuxer) run() {
} }
}() }()
r.ctxCancel() m.ctxCancel()
for _, req := range r.requests { for _, req := range m.requests {
req.Res <- hlsMuxerResponse{Status: http.StatusNotFound} req.Res <- hlsMuxerResponse{Status: http.StatusNotFound}
} }
r.parent.onMuxerClose(r) m.parent.onMuxerClose(m)
r.log(logger.Info, "closed (%v)", err) m.log(logger.Info, "closed (%v)", err)
} }
func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{}) error { func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{}) error {
res := r.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{ res := m.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{
Author: r, Author: m,
PathName: r.pathName, PathName: m.pathName,
IP: nil, IP: nil,
ValidateCredentials: nil, ValidateCredentials: nil,
}) })
@ -257,10 +257,10 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
return res.Err return res.Err
} }
r.path = res.Path m.path = res.Path
defer func() { defer func() {
r.path.onReaderRemove(pathReaderRemoveReq{Author: r}) m.path.onReaderRemove(pathReaderRemoveReq{Author: m})
}() }()
var videoTrack *gortsplib.Track var videoTrack *gortsplib.Track
@ -302,28 +302,28 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
} }
var err error var err error
r.muxer, err = hls.NewMuxer( m.muxer, err = hls.NewMuxer(
r.hlsSegmentCount, m.hlsSegmentCount,
time.Duration(r.hlsSegmentDuration), time.Duration(m.hlsSegmentDuration),
videoTrack, videoTrack,
audioTrack, audioTrack,
) )
if err != nil { if err != nil {
return err return err
} }
defer r.muxer.Close() defer m.muxer.Close()
innerReady <- struct{}{} innerReady <- struct{}{}
r.ringBuffer = ringbuffer.New(uint64(r.readBufferCount)) m.ringBuffer = ringbuffer.New(uint64(m.readBufferCount))
r.path.onReaderPlay(pathReaderPlayReq{Author: r}) m.path.onReaderPlay(pathReaderPlayReq{Author: m})
writerDone := make(chan error) writerDone := make(chan error)
go func() { go func() {
writerDone <- func() error { writerDone <- func() error {
for { for {
data, ok := r.ringBuffer.Pull() data, ok := m.ringBuffer.Pull()
if !ok { if !ok {
return fmt.Errorf("terminated") return fmt.Errorf("terminated")
} }
@ -333,7 +333,7 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
var pkt rtp.Packet var pkt rtp.Packet
err := pkt.Unmarshal(pair.buf) err := pkt.Unmarshal(pair.buf)
if err != nil { if err != nil {
r.log(logger.Warn, "unable to decode RTP packet: %v", err) m.log(logger.Warn, "unable to decode RTP packet: %v", err)
continue continue
} }
@ -341,12 +341,12 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
if err != nil { if err != nil {
if err != rtph264.ErrMorePacketsNeeded && if err != rtph264.ErrMorePacketsNeeded &&
err != rtph264.ErrNonStartingPacketAndNoPrevious { err != rtph264.ErrNonStartingPacketAndNoPrevious {
r.log(logger.Warn, "unable to decode video track: %v", err) m.log(logger.Warn, "unable to decode video track: %v", err)
} }
continue continue
} }
err = r.muxer.WriteH264(pts, nalus) err = m.muxer.WriteH264(pts, nalus)
if err != nil { if err != nil {
return err return err
} }
@ -354,19 +354,19 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
var pkt rtp.Packet var pkt rtp.Packet
err := pkt.Unmarshal(pair.buf) err := pkt.Unmarshal(pair.buf)
if err != nil { if err != nil {
r.log(logger.Warn, "unable to decode RTP packet: %v", err) m.log(logger.Warn, "unable to decode RTP packet: %v", err)
continue continue
} }
aus, pts, err := aacDecoder.Decode(&pkt) aus, pts, err := aacDecoder.Decode(&pkt)
if err != nil { if err != nil {
if err != rtpaac.ErrMorePacketsNeeded { if err != rtpaac.ErrMorePacketsNeeded {
r.log(logger.Warn, "unable to decode audio track: %v", err) m.log(logger.Warn, "unable to decode audio track: %v", err)
} }
continue continue
} }
err = r.muxer.WriteAAC(pts, aus) err = m.muxer.WriteAAC(pts, aus)
if err != nil { if err != nil {
return err return err
} }
@ -381,9 +381,9 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
for { for {
select { select {
case <-closeCheckTicker.C: case <-closeCheckTicker.C:
t := time.Unix(atomic.LoadInt64(r.lastRequestTime), 0) t := time.Unix(atomic.LoadInt64(m.lastRequestTime), 0)
if !r.hlsAlwaysRemux && time.Since(t) >= closeAfterInactivity { if !m.hlsAlwaysRemux && time.Since(t) >= closeAfterInactivity {
r.ringBuffer.Close() m.ringBuffer.Close()
<-writerDone <-writerDone
return nil return nil
} }
@ -392,23 +392,23 @@ func (r *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
return err return err
case <-innerCtx.Done(): case <-innerCtx.Done():
r.ringBuffer.Close() m.ringBuffer.Close()
<-writerDone <-writerDone
return nil return nil
} }
} }
} }
func (r *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse { func (m *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse {
atomic.StoreInt64(r.lastRequestTime, time.Now().Unix()) atomic.StoreInt64(m.lastRequestTime, time.Now().Unix())
conf := r.path.Conf() conf := m.path.Conf()
if conf.ReadIPs != nil { if conf.ReadIPs != nil {
tmp, _, _ := net.SplitHostPort(req.Req.RemoteAddr) tmp, _, _ := net.SplitHostPort(req.Req.RemoteAddr)
ip := net.ParseIP(tmp) ip := net.ParseIP(tmp)
if !ipEqualOrInRange(ip, conf.ReadIPs) { if !ipEqualOrInRange(ip, conf.ReadIPs) {
r.log(logger.Info, "ip '%s' not allowed", ip) m.log(logger.Info, "ip '%s' not allowed", ip)
return hlsMuxerResponse{Status: http.StatusUnauthorized} return hlsMuxerResponse{Status: http.StatusUnauthorized}
} }
} }
@ -432,7 +432,7 @@ func (r *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse {
Header: map[string]string{ Header: map[string]string{
"Content-Type": `application/x-mpegURL`, "Content-Type": `application/x-mpegURL`,
}, },
Body: r.muxer.PrimaryPlaylist(), Body: m.muxer.PrimaryPlaylist(),
} }
case req.File == "stream.m3u8": case req.File == "stream.m3u8":
@ -441,11 +441,11 @@ func (r *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse {
Header: map[string]string{ Header: map[string]string{
"Content-Type": `application/x-mpegURL`, "Content-Type": `application/x-mpegURL`,
}, },
Body: r.muxer.StreamPlaylist(), Body: m.muxer.StreamPlaylist(),
} }
case strings.HasSuffix(req.File, ".ts"): case strings.HasSuffix(req.File, ".ts"):
r := r.muxer.Segment(req.File) r := m.muxer.Segment(req.File)
if r == nil { if r == nil {
return hlsMuxerResponse{Status: http.StatusNotFound} return hlsMuxerResponse{Status: http.StatusNotFound}
} }
@ -473,28 +473,28 @@ func (r *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse {
} }
// onRequest is called by hlsserver.Server (forwarded from ServeHTTP). // onRequest is called by hlsserver.Server (forwarded from ServeHTTP).
func (r *hlsMuxer) onRequest(req hlsMuxerRequest) { func (m *hlsMuxer) onRequest(req hlsMuxerRequest) {
select { select {
case r.request <- req: case m.request <- req:
case <-r.ctx.Done(): case <-m.ctx.Done():
req.Res <- hlsMuxerResponse{Status: http.StatusNotFound} req.Res <- hlsMuxerResponse{Status: http.StatusNotFound}
} }
} }
// onReaderAccepted implements reader. // onReaderAccepted implements reader.
func (r *hlsMuxer) onReaderAccepted() { func (m *hlsMuxer) onReaderAccepted() {
r.log(logger.Info, "is converting into HLS") m.log(logger.Info, "is converting into HLS")
} }
// onReaderFrame implements reader. // onReaderFrame implements reader.
func (r *hlsMuxer) onReaderFrame(trackID int, streamType gortsplib.StreamType, payload []byte) { func (m *hlsMuxer) onReaderFrame(trackID int, streamType gortsplib.StreamType, payload []byte) {
if streamType == gortsplib.StreamTypeRTP { if streamType == gortsplib.StreamTypeRTP {
r.ringBuffer.Push(hlsMuxerTrackIDPayloadPair{trackID, payload}) m.ringBuffer.Push(hlsMuxerTrackIDPayloadPair{trackID, payload})
} }
} }
// onReaderAPIDescribe implements reader. // onReaderAPIDescribe implements reader.
func (r *hlsMuxer) onReaderAPIDescribe() interface{} { func (m *hlsMuxer) onReaderAPIDescribe() interface{} {
return struct { return struct {
Type string `json:"type"` Type string `json:"type"`
}{"hlsMuxer"} }{"hlsMuxer"}

Loading…
Cancel
Save