Browse Source

hls: fix freeze when sourceOnDemand is yes and multiple sources are requested at the same time (#493)

pull/509/head
aler9 5 years ago
parent
commit
3872b42434
  1. 170
      internal/core/hls_remuxer.go

170
internal/core/hls_remuxer.go

@ -111,6 +111,7 @@ type hlsRemuxer struct {
ringBuffer *ringbuffer.RingBuffer ringBuffer *ringbuffer.RingBuffer
lastRequestTime *int64 lastRequestTime *int64
muxer *hls.Muxer muxer *hls.Muxer
requests []hlsRemuxerRequest
// in // in
request chan hlsRemuxerRequest request chan hlsRemuxerRequest
@ -177,22 +178,44 @@ func (r *hlsRemuxer) PathName() string {
func (r *hlsRemuxer) run() { func (r *hlsRemuxer) run() {
defer r.wg.Done() defer r.wg.Done()
innerCtx, innerCtxCancel := context.WithCancel(context.Background()) remuxerCtx, remuxerCtxCancel := context.WithCancel(context.Background())
runErr := make(chan error) remuxerReady := make(chan struct{})
remuxerErr := make(chan error)
go func() { go func() {
runErr <- r.runInner(innerCtx) remuxerErr <- r.runRemuxer(remuxerCtx, remuxerReady)
}() }()
select { isReady := false
case err := <-runErr:
innerCtxCancel()
if err != nil {
r.log(logger.Info, "ERR: %s", err)
}
case <-r.ctx.Done(): outer:
innerCtxCancel() for {
<-runErr select {
case <-r.ctx.Done():
remuxerCtxCancel()
<-remuxerErr
break outer
case req := <-r.request:
if isReady {
r.handleRequest(req)
} else {
r.requests = append(r.requests, req)
}
case <-remuxerReady:
isReady = true
for _, req := range r.requests {
r.handleRequest(req)
}
r.requests = nil
case err := <-remuxerErr:
remuxerCtxCancel()
if err != nil {
r.log(logger.Info, "ERR: %s", err)
}
break outer
}
} }
r.ctxCancel() r.ctxCancel()
@ -200,14 +223,13 @@ func (r *hlsRemuxer) run() {
r.parent.OnRemuxerClose(r) r.parent.OnRemuxerClose(r)
} }
func (r *hlsRemuxer) runInner(innerCtx context.Context) error { func (r *hlsRemuxer) runRemuxer(remuxerCtx context.Context, remuxerReady chan struct{}) error {
res := r.pathManager.OnReaderSetupPlay(pathReaderSetupPlayReq{ res := r.pathManager.OnReaderSetupPlay(pathReaderSetupPlayReq{
Author: r, Author: r,
PathName: r.pathName, PathName: r.pathName,
IP: nil, IP: nil,
ValidateCredentials: nil, ValidateCredentials: nil,
}) })
if res.Err != nil { if res.Err != nil {
return res.Err return res.Err
} }
@ -283,21 +305,11 @@ func (r *hlsRemuxer) runInner(innerCtx context.Context) error {
} }
defer r.muxer.Close() defer r.muxer.Close()
// start request handler only after muxer has been inizialized remuxerReady <- struct{}{}
requestHandlerTerminate := make(chan struct{})
requestHandlerDone := make(chan struct{})
go r.runRequestHandler(requestHandlerTerminate, requestHandlerDone)
defer func() {
close(requestHandlerTerminate)
<-requestHandlerDone
}()
r.ringBuffer = ringbuffer.New(uint64(r.readBufferCount)) r.ringBuffer = ringbuffer.New(uint64(r.readBufferCount))
r.path.OnReaderPlay(pathReaderPlayReq{ r.path.OnReaderPlay(pathReaderPlayReq{Author: r})
Author: r,
})
writerDone := make(chan error) writerDone := make(chan error)
go func() { go func() {
@ -396,7 +408,7 @@ func (r *hlsRemuxer) runInner(innerCtx context.Context) error {
case err := <-writerDone: case err := <-writerDone:
return err return err
case <-innerCtx.Done(): case <-remuxerCtx.Done():
r.ringBuffer.Close() r.ringBuffer.Close()
<-writerDone <-writerDone
return nil return nil
@ -404,73 +416,61 @@ func (r *hlsRemuxer) runInner(innerCtx context.Context) error {
} }
} }
func (r *hlsRemuxer) runRequestHandler(terminate chan struct{}, done chan struct{}) { func (r *hlsRemuxer) handleRequest(req hlsRemuxerRequest) {
defer close(done) atomic.StoreInt64(r.lastRequestTime, time.Now().Unix())
for {
select {
case <-terminate:
return
case preq := <-r.request:
req := preq
atomic.StoreInt64(r.lastRequestTime, time.Now().Unix()) conf := r.path.Conf()
conf := r.path.Conf() if conf.ReadIPsParsed != nil {
tmp, _, _ := net.SplitHostPort(req.Req.RemoteAddr)
if conf.ReadIPsParsed != nil { ip := net.ParseIP(tmp)
tmp, _, _ := net.SplitHostPort(req.Req.RemoteAddr) if !ipEqualOrInRange(ip, conf.ReadIPsParsed) {
ip := net.ParseIP(tmp) r.log(logger.Info, "ERR: ip '%s' not allowed", ip)
if !ipEqualOrInRange(ip, conf.ReadIPsParsed) { req.W.WriteHeader(http.StatusUnauthorized)
r.log(logger.Info, "ERR: ip '%s' not allowed", ip) req.Res <- nil
req.W.WriteHeader(http.StatusUnauthorized) return
req.Res <- nil }
continue }
}
}
if conf.ReadUser != "" { if conf.ReadUser != "" {
user, pass, ok := req.Req.BasicAuth() user, pass, ok := req.Req.BasicAuth()
if !ok || user != conf.ReadUser || pass != conf.ReadPass { if !ok || user != conf.ReadUser || pass != conf.ReadPass {
req.W.Header().Set("WWW-Authenticate", `Basic realm="rtsp-simple-server"`) req.W.Header().Set("WWW-Authenticate", `Basic realm="rtsp-simple-server"`)
req.W.WriteHeader(http.StatusUnauthorized) req.W.WriteHeader(http.StatusUnauthorized)
req.Res <- nil req.Res <- nil
continue return
} }
} }
switch { switch {
case req.File == "stream.m3u8": case req.File == "stream.m3u8":
r := r.muxer.Playlist() r := r.muxer.Playlist()
if r == nil { if r == nil {
req.W.WriteHeader(http.StatusNotFound) req.W.WriteHeader(http.StatusNotFound)
req.Res <- nil req.Res <- nil
continue return
} }
req.W.Header().Set("Content-Type", `application/x-mpegURL`) req.W.Header().Set("Content-Type", `application/x-mpegURL`)
req.Res <- r req.Res <- r
case strings.HasSuffix(req.File, ".ts"): case strings.HasSuffix(req.File, ".ts"):
r := r.muxer.TSFile(req.File) r := r.muxer.TSFile(req.File)
if r == nil { if r == nil {
req.W.WriteHeader(http.StatusNotFound) req.W.WriteHeader(http.StatusNotFound)
req.Res <- nil req.Res <- nil
continue return
} }
req.W.Header().Set("Content-Type", `video/MP2T`) req.W.Header().Set("Content-Type", `video/MP2T`)
req.Res <- r req.Res <- r
case req.File == "": case req.File == "":
req.Res <- bytes.NewReader([]byte(index)) req.Res <- bytes.NewReader([]byte(index))
default: default:
req.W.WriteHeader(http.StatusNotFound) req.W.WriteHeader(http.StatusNotFound)
req.Res <- nil req.Res <- nil
}
}
} }
} }

Loading…
Cancel
Save