mediamtx/internal/core/path.go

979 lines
23 KiB
Go
Raw Normal View History

package core
2020-10-19 20:17:48 +00:00
import (
2021-05-10 19:32:59 +00:00
"context"
2020-10-19 20:17:48 +00:00
"fmt"
"net"
2020-10-19 20:17:48 +00:00
"strings"
"sync"
"time"
"github.com/aler9/gortsplib"
2021-02-09 21:33:50 +00:00
"github.com/aler9/gortsplib/pkg/base"
2020-10-19 20:17:48 +00:00
2020-11-01 21:56:56 +00:00
"github.com/aler9/rtsp-simple-server/internal/conf"
"github.com/aler9/rtsp-simple-server/internal/externalcmd"
"github.com/aler9/rtsp-simple-server/internal/logger"
2020-10-19 20:17:48 +00:00
)
func newEmptyTimer() *time.Timer {
t := time.NewTimer(0)
<-t.C
return t
}
2020-10-19 20:17:48 +00:00
type pathErrNoOnePublishing struct {
PathName string
}
// Error implements the error interface.
func (e pathErrNoOnePublishing) Error() string {
return fmt.Sprintf("no one is publishing to path '%s'", e.PathName)
}
type pathErrAuthNotCritical struct {
*base.Response
}
// Error implements the error interface.
func (pathErrAuthNotCritical) Error() string {
return "non-critical authentication error"
}
type pathErrAuthCritical struct {
Message string
Response *base.Response
}
// Error implements the error interface.
func (pathErrAuthCritical) Error() string {
return "critical authentication error"
}
type pathParent interface {
2021-10-27 19:01:00 +00:00
log(logger.Level, string, ...interface{})
onPathSourceReady(*path)
onPathClose(*path)
2020-10-19 20:17:48 +00:00
}
type pathRTSPSession interface {
IsRTSPSession()
}
2020-10-27 23:29:53 +00:00
type sourceRedirect struct{}
2021-10-27 19:01:00 +00:00
// onSourceAPIDescribe implements source.
func (*sourceRedirect) onSourceAPIDescribe() interface{} {
return struct {
Type string `json:"type"`
}{"redirect"}
}
2020-10-27 23:29:53 +00:00
type pathReaderState int
2020-10-19 20:17:48 +00:00
const (
pathReaderStatePrePlay pathReaderState = iota
pathReaderStatePlay
2020-10-19 20:17:48 +00:00
)
2021-08-03 20:40:47 +00:00
type pathOnDemandState int
const (
2021-08-03 20:40:47 +00:00
pathOnDemandStateInitial pathOnDemandState = iota
pathOnDemandStateWaitingReady
pathOnDemandStateReady
pathOnDemandStateClosing
)
type pathSourceStaticSetReadyRes struct {
Stream *stream
Err error
}
type pathSourceStaticSetReadyReq struct {
Source sourceStatic
Tracks gortsplib.Tracks
Res chan pathSourceStaticSetReadyRes
}
type pathSourceStaticSetNotReadyReq struct {
2021-08-03 20:40:47 +00:00
Source sourceStatic
Res chan struct{}
}
type pathReaderRemoveReq struct {
Author reader
Res chan struct{}
}
type pathPublisherRemoveReq struct {
Author publisher
Res chan struct{}
}
type pathDescribeRes struct {
2021-08-01 15:22:28 +00:00
Path *path
Stream *stream
Redirect string
Err error
}
type pathDescribeReq struct {
PathName string
URL *base.URL
IP net.IP
ValidateCredentials func(pathUser conf.Credential, pathPass conf.Credential) error
Res chan pathDescribeRes
}
type pathReaderSetupPlayRes struct {
Path *path
Stream *stream
Err error
}
type pathReaderSetupPlayReq struct {
Author reader
PathName string
IP net.IP
ValidateCredentials func(pathUser conf.Credential, pathPass conf.Credential) error
Res chan pathReaderSetupPlayRes
}
type pathPublisherAnnounceRes struct {
Path *path
Err error
}
type pathPublisherAnnounceReq struct {
Author publisher
PathName string
IP net.IP
ValidateCredentials func(pathUser conf.Credential, pathPass conf.Credential) error
Res chan pathPublisherAnnounceRes
}
type pathReaderPlayReq struct {
Author reader
Res chan struct{}
}
type pathPublisherRecordRes struct {
Stream *stream
Err error
}
type pathPublisherRecordReq struct {
Author publisher
Tracks gortsplib.Tracks
Res chan pathPublisherRecordRes
}
type pathReaderPauseReq struct {
Author reader
Res chan struct{}
}
type pathPublisherPauseReq struct {
Author publisher
Res chan struct{}
}
type path struct {
rtspAddress string
readTimeout conf.StringDuration
writeTimeout conf.StringDuration
readBufferCount int
readBufferSize int
2021-01-10 11:55:53 +00:00
confName string
conf *conf.PathConf
name string
wg *sync.WaitGroup
parent pathParent
2020-10-19 20:17:48 +00:00
2021-08-03 20:40:47 +00:00
ctx context.Context
ctxCancel func()
source source
sourceReady bool
sourceStaticWg sync.WaitGroup
readers map[reader]pathReaderState
describeRequests []pathDescribeReq
setupPlayRequests []pathReaderSetupPlayReq
stream *stream
2021-08-03 20:40:47 +00:00
onDemandCmd *externalcmd.Cmd
onPublishCmd *externalcmd.Cmd
onDemandReadyTimer *time.Timer
onDemandCloseTimer *time.Timer
onDemandState pathOnDemandState
2020-10-19 20:17:48 +00:00
// in
sourceStaticSetReady chan pathSourceStaticSetReadyReq
sourceStaticSetNotReady chan pathSourceStaticSetNotReadyReq
describe chan pathDescribeReq
publisherRemove chan pathPublisherRemoveReq
publisherAnnounce chan pathPublisherAnnounceReq
publisherRecord chan pathPublisherRecordReq
publisherPause chan pathPublisherPauseReq
readerRemove chan pathReaderRemoveReq
readerSetupPlay chan pathReaderSetupPlayReq
readerPlay chan pathReaderPlayReq
readerPause chan pathReaderPauseReq
2021-11-05 16:14:31 +00:00
apiPathsList chan apiPathsListSubReq
}
func newPath(
parentCtx context.Context,
rtspAddress string,
readTimeout conf.StringDuration,
writeTimeout conf.StringDuration,
readBufferCount int,
readBufferSize int,
confName string,
2020-10-19 20:17:48 +00:00
conf *conf.PathConf,
name string,
wg *sync.WaitGroup,
parent pathParent) *path {
ctx, ctxCancel := context.WithCancel(parentCtx)
2021-05-10 19:32:59 +00:00
pa := &path{
rtspAddress: rtspAddress,
readTimeout: readTimeout,
writeTimeout: writeTimeout,
readBufferCount: readBufferCount,
readBufferSize: readBufferSize,
confName: confName,
conf: conf,
name: name,
wg: wg,
parent: parent,
ctx: ctx,
ctxCancel: ctxCancel,
readers: make(map[reader]pathReaderState),
2021-08-03 20:40:47 +00:00
onDemandReadyTimer: newEmptyTimer(),
onDemandCloseTimer: newEmptyTimer(),
sourceStaticSetReady: make(chan pathSourceStaticSetReadyReq),
sourceStaticSetNotReady: make(chan pathSourceStaticSetNotReadyReq),
describe: make(chan pathDescribeReq),
publisherRemove: make(chan pathPublisherRemoveReq),
publisherAnnounce: make(chan pathPublisherAnnounceReq),
publisherRecord: make(chan pathPublisherRecordReq),
publisherPause: make(chan pathPublisherPauseReq),
readerRemove: make(chan pathReaderRemoveReq),
readerSetupPlay: make(chan pathReaderSetupPlayReq),
readerPlay: make(chan pathReaderPlayReq),
readerPause: make(chan pathReaderPauseReq),
2021-11-05 16:14:31 +00:00
apiPathsList: make(chan apiPathsListSubReq),
2020-10-19 20:17:48 +00:00
}
pa.log(logger.Info, "opened")
2020-10-19 20:17:48 +00:00
pa.wg.Add(1)
go pa.run()
2021-08-01 14:58:46 +00:00
2020-10-19 20:17:48 +00:00
return pa
}
2021-10-27 19:01:00 +00:00
func (pa *path) close() {
2021-05-10 19:32:59 +00:00
pa.ctxCancel()
2020-10-19 20:17:48 +00:00
}
2020-11-05 11:30:25 +00:00
// Log is the main logging function.
2021-10-27 19:01:00 +00:00
func (pa *path) log(level logger.Level, format string, args ...interface{}) {
pa.parent.log(level, "[path "+pa.name+"] "+format, args...)
2020-10-19 20:17:48 +00:00
}
2021-07-31 16:32:00 +00:00
// ConfName returns the configuration name of this path.
func (pa *path) ConfName() string {
return pa.confName
}
// Conf returns the configuration of this path.
func (pa *path) Conf() *conf.PathConf {
return pa.conf
}
// Name returns the name of this path.
func (pa *path) Name() string {
return pa.name
}
func (pa *path) run() {
2020-10-19 20:17:48 +00:00
defer pa.wg.Done()
if pa.conf.Source == "redirect" {
2020-10-27 23:29:53 +00:00
pa.source = &sourceRedirect{}
2021-08-03 20:40:47 +00:00
} else if !pa.conf.SourceOnDemand && pa.hasStaticSource() {
pa.staticSourceCreate()
2020-10-19 20:17:48 +00:00
}
2020-11-21 12:35:33 +00:00
var onInitCmd *externalcmd.Cmd
2020-10-19 20:17:48 +00:00
if pa.conf.RunOnInit != "" {
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnInit command started")
_, port, _ := net.SplitHostPort(pa.rtspAddress)
onInitCmd = externalcmd.New(pa.conf.RunOnInit, pa.conf.RunOnInitRestart, externalcmd.Environment{
2020-11-01 16:33:06 +00:00
Path: pa.name,
Port: port,
2020-11-01 16:33:06 +00:00
})
2020-10-19 20:17:48 +00:00
}
err := func() error {
for {
select {
case <-pa.onDemandReadyTimer.C:
for _, req := range pa.describeRequests {
req.Res <- pathDescribeRes{Err: fmt.Errorf("source of path '%s' has timed out", pa.name)}
}
pa.describeRequests = nil
2020-10-19 20:17:48 +00:00
for _, req := range pa.setupPlayRequests {
req.Res <- pathReaderSetupPlayRes{Err: fmt.Errorf("source of path '%s' has timed out", pa.name)}
}
pa.setupPlayRequests = nil
pa.onDemandCloseSource()
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
case <-pa.onDemandCloseTimer.C:
pa.onDemandCloseSource()
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
2020-10-19 20:17:48 +00:00
case req := <-pa.sourceStaticSetReady:
if req.Source == pa.source {
pa.sourceSetReady(req.Tracks)
req.Res <- pathSourceStaticSetReadyRes{Stream: pa.stream}
} else {
req.Res <- pathSourceStaticSetReadyRes{Err: fmt.Errorf("terminated")}
}
2020-10-19 20:17:48 +00:00
case req := <-pa.sourceStaticSetNotReady:
if req.Source == pa.source {
if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial {
pa.onDemandCloseSource()
} else {
pa.sourceSetNotReady()
}
}
close(req.Res)
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
2021-08-03 20:40:47 +00:00
case req := <-pa.describe:
pa.handleDescribe(req)
2020-10-19 20:17:48 +00:00
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
case req := <-pa.publisherRemove:
pa.handlePublisherRemove(req)
2020-10-19 20:17:48 +00:00
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
2021-08-03 20:40:47 +00:00
case req := <-pa.publisherAnnounce:
pa.handlePublisherAnnounce(req)
2021-03-10 14:06:45 +00:00
case req := <-pa.publisherRecord:
pa.handlePublisherRecord(req)
2020-10-19 20:17:48 +00:00
case req := <-pa.publisherPause:
pa.handlePublisherPause(req)
2020-10-19 20:17:48 +00:00
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
2021-08-03 20:40:47 +00:00
case req := <-pa.readerRemove:
pa.handleReaderRemove(req)
case req := <-pa.readerSetupPlay:
pa.handleReaderSetupPlay(req)
if pa.shouldClose() {
return fmt.Errorf("not in use")
}
case req := <-pa.readerPlay:
pa.handleReaderPlay(req)
case req := <-pa.readerPause:
pa.handleReaderPause(req)
2020-10-19 20:17:48 +00:00
case req := <-pa.apiPathsList:
pa.handleAPIPathsList(req)
case <-pa.ctx.Done():
return fmt.Errorf("terminated")
}
2020-10-19 20:17:48 +00:00
}
}()
2020-10-19 20:17:48 +00:00
2021-05-10 19:32:59 +00:00
pa.ctxCancel()
2021-08-03 20:40:47 +00:00
pa.onDemandReadyTimer.Stop()
pa.onDemandCloseTimer.Stop()
if onInitCmd != nil {
onInitCmd.Close()
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnInit command stopped")
}
for _, req := range pa.describeRequests {
req.Res <- pathDescribeRes{Err: fmt.Errorf("terminated")}
}
for _, req := range pa.setupPlayRequests {
req.Res <- pathReaderSetupPlayRes{Err: fmt.Errorf("terminated")}
}
for rp := range pa.readers {
2021-10-27 19:01:00 +00:00
rp.close()
}
if pa.stream != nil {
pa.stream.close()
}
if pa.source != nil {
if source, ok := pa.source.(sourceStatic); ok {
2021-10-27 19:01:00 +00:00
source.close()
pa.sourceStaticWg.Wait()
} else if source, ok := pa.source.(publisher); ok {
2021-10-27 19:01:00 +00:00
source.close()
}
}
2020-10-19 20:17:48 +00:00
// close onDemandCmd after the source has been closed.
// this avoids a deadlock in which onDemandCmd is a
// RTSP publisher that sends a TEARDOWN request and waits
// for the response (like FFmpeg), but it can't since
2021-10-03 14:08:10 +00:00
// the path is already waiting for the command to close.
if pa.onDemandCmd != nil {
pa.onDemandCmd.Close()
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnDemand command stopped")
}
pa.log(logger.Info, "closed (%v)", err)
2021-10-27 19:01:00 +00:00
pa.parent.onPathClose(pa)
2020-10-19 20:17:48 +00:00
}
func (pa *path) shouldClose() bool {
return pa.conf.Regexp != nil &&
pa.source == nil &&
len(pa.readers) == 0 &&
len(pa.describeRequests) == 0 &&
len(pa.setupPlayRequests) == 0
}
func (pa *path) hasStaticSource() bool {
return strings.HasPrefix(pa.conf.Source, "rtsp://") ||
2020-12-14 22:32:24 +00:00
strings.HasPrefix(pa.conf.Source, "rtsps://") ||
2021-09-05 15:35:04 +00:00
strings.HasPrefix(pa.conf.Source, "rtmp://") ||
strings.HasPrefix(pa.conf.Source, "http://") ||
strings.HasPrefix(pa.conf.Source, "https://")
}
2021-08-03 20:40:47 +00:00
func (pa *path) isOnDemand() bool {
return (pa.hasStaticSource() && pa.conf.SourceOnDemand) || pa.conf.RunOnDemand != ""
}
func (pa *path) onDemandStartSource() {
pa.onDemandReadyTimer.Stop()
if pa.hasStaticSource() {
pa.staticSourceCreate()
pa.onDemandReadyTimer = time.NewTimer(time.Duration(pa.conf.SourceOnDemandStartTimeout))
2021-08-03 20:40:47 +00:00
} else {
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnDemand command started")
2021-08-03 20:40:47 +00:00
_, port, _ := net.SplitHostPort(pa.rtspAddress)
pa.onDemandCmd = externalcmd.New(pa.conf.RunOnDemand, pa.conf.RunOnDemandRestart, externalcmd.Environment{
Path: pa.name,
Port: port,
})
pa.onDemandReadyTimer = time.NewTimer(time.Duration(pa.conf.RunOnDemandStartTimeout))
2021-08-03 20:40:47 +00:00
}
pa.onDemandState = pathOnDemandStateWaitingReady
}
func (pa *path) onDemandScheduleClose() {
pa.onDemandCloseTimer.Stop()
if pa.hasStaticSource() {
pa.onDemandCloseTimer = time.NewTimer(time.Duration(pa.conf.SourceOnDemandCloseAfter))
2021-08-03 20:40:47 +00:00
} else {
pa.onDemandCloseTimer = time.NewTimer(time.Duration(pa.conf.RunOnDemandCloseAfter))
2021-08-03 20:40:47 +00:00
}
pa.onDemandState = pathOnDemandStateClosing
}
func (pa *path) onDemandCloseSource() {
if pa.onDemandState == pathOnDemandStateClosing {
pa.onDemandCloseTimer.Stop()
pa.onDemandCloseTimer = newEmptyTimer()
}
// set state before doPublisherRemove()
pa.onDemandState = pathOnDemandStateInitial
if pa.hasStaticSource() {
if pa.sourceReady {
pa.sourceSetNotReady()
}
2021-10-27 19:01:00 +00:00
pa.source.(sourceStatic).close()
pa.source = nil
2021-08-03 20:40:47 +00:00
} else {
if pa.source != nil {
2021-10-27 19:01:00 +00:00
pa.source.(publisher).close()
pa.doPublisherRemove()
}
// close onDemandCmd after the source has been closed.
// this avoids a deadlock in which onDemandCmd is a
// RTSP publisher that sends a TEARDOWN request and waits
// for the response (like FFmpeg), but it can't since
2021-10-03 14:08:10 +00:00
// the path is already waiting for the command to close.
if pa.onDemandCmd != nil {
pa.onDemandCmd.Close()
pa.onDemandCmd = nil
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnDemand command stopped")
}
2021-08-03 20:40:47 +00:00
}
}
func (pa *path) sourceSetReady(tracks gortsplib.Tracks) {
2021-08-03 20:40:47 +00:00
pa.sourceReady = true
pa.stream = newStream(tracks)
2021-08-03 20:40:47 +00:00
if pa.isOnDemand() {
pa.onDemandReadyTimer.Stop()
pa.onDemandReadyTimer = newEmptyTimer()
for _, req := range pa.describeRequests {
req.Res <- pathDescribeRes{
Stream: pa.stream,
}
}
pa.describeRequests = nil
for _, req := range pa.setupPlayRequests {
2021-08-11 10:45:53 +00:00
pa.handleReaderSetupPlayPost(req)
2021-08-03 20:40:47 +00:00
}
pa.setupPlayRequests = nil
if len(pa.readers) > 0 {
pa.onDemandState = pathOnDemandStateReady
} else {
pa.onDemandScheduleClose()
}
}
2021-10-27 19:01:00 +00:00
pa.parent.onPathSourceReady(pa)
2021-08-03 20:40:47 +00:00
}
func (pa *path) sourceSetNotReady() {
for r := range pa.readers {
pa.doReaderRemove(r)
2021-10-27 19:01:00 +00:00
r.close()
}
// close onPublishCmd after all readers have been closed.
// this avoids a deadlock in which onPublishCmd is a
// RTSP reader that sends a TEARDOWN request and waits
// for the response (like FFmpeg), but it can't since
2021-10-03 14:08:10 +00:00
// the path is already waiting for the command to close.
2021-08-03 20:40:47 +00:00
if pa.onPublishCmd != nil {
pa.onPublishCmd.Close()
pa.onPublishCmd = nil
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnPublish command stopped")
2021-08-03 20:40:47 +00:00
}
pa.sourceReady = false
pa.stream.close()
pa.stream = nil
2021-08-03 20:40:47 +00:00
}
func (pa *path) staticSourceCreate() {
2021-09-05 15:35:04 +00:00
switch {
case strings.HasPrefix(pa.conf.Source, "rtsp://") ||
strings.HasPrefix(pa.conf.Source, "rtsps://"):
pa.source = newRTSPSource(
2021-05-11 15:20:32 +00:00
pa.ctx,
2021-01-10 11:55:53 +00:00
pa.conf.Source,
pa.conf.SourceProtocol,
pa.conf.SourceAnyPortEnable,
pa.conf.SourceFingerprint,
2021-01-10 11:55:53 +00:00
pa.readTimeout,
pa.writeTimeout,
pa.readBufferCount,
pa.readBufferSize,
&pa.sourceStaticWg,
2021-01-10 11:55:53 +00:00
pa)
2021-09-05 15:35:04 +00:00
case strings.HasPrefix(pa.conf.Source, "rtmp://"):
pa.source = newRTMPSource(
2021-05-11 15:20:32 +00:00
pa.ctx,
2021-01-10 11:55:53 +00:00
pa.conf.Source,
2021-01-31 15:24:58 +00:00
pa.readTimeout,
pa.writeTimeout,
&pa.sourceStaticWg,
2021-01-10 11:55:53 +00:00
pa)
2021-09-05 15:35:04 +00:00
case strings.HasPrefix(pa.conf.Source, "http://") ||
strings.HasPrefix(pa.conf.Source, "https://"):
pa.source = newHLSSource(
pa.ctx,
pa.conf.Source,
pa.conf.SourceFingerprint,
2021-09-05 15:35:04 +00:00
&pa.sourceStaticWg,
pa)
}
2020-10-19 20:17:48 +00:00
}
2021-08-03 20:40:47 +00:00
func (pa *path) doReaderRemove(r reader) {
state := pa.readers[r]
if state == pathReaderStatePlay {
pa.stream.readerRemove(r)
2020-10-19 20:17:48 +00:00
}
delete(pa.readers, r)
2020-10-19 20:17:48 +00:00
}
2021-08-03 20:40:47 +00:00
func (pa *path) doPublisherRemove() {
if pa.sourceReady {
if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial {
pa.onDemandCloseSource()
} else {
pa.sourceSetNotReady()
}
2021-10-03 14:01:40 +00:00
} else {
for r := range pa.readers {
pa.doReaderRemove(r)
2021-10-27 19:01:00 +00:00
r.close()
2021-10-03 14:01:40 +00:00
}
}
pa.source = nil
}
2021-08-11 10:45:53 +00:00
func (pa *path) handleDescribe(req pathDescribeReq) {
if _, ok := pa.source.(*sourceRedirect); ok {
req.Res <- pathDescribeRes{
Redirect: pa.conf.SourceRedirect,
}
return
}
2021-08-03 20:40:47 +00:00
if pa.sourceReady {
req.Res <- pathDescribeRes{
Stream: pa.stream,
}
return
2021-08-03 20:40:47 +00:00
}
2021-08-03 20:40:47 +00:00
if pa.isOnDemand() {
if pa.onDemandState == pathOnDemandStateInitial {
pa.onDemandStartSource()
}
pa.describeRequests = append(pa.describeRequests, req)
return
2021-08-03 20:40:47 +00:00
}
2021-08-03 20:40:47 +00:00
if pa.conf.Fallback != "" {
fallbackURL := func() string {
if strings.HasPrefix(pa.conf.Fallback, "/") {
ur := base.URL{
Scheme: req.URL.Scheme,
User: req.URL.User,
Host: req.URL.Host,
Path: pa.conf.Fallback,
2021-02-09 21:33:50 +00:00
}
2021-08-03 20:40:47 +00:00
return ur.String()
}
return pa.conf.Fallback
}()
req.Res <- pathDescribeRes{Redirect: fallbackURL}
return
2020-10-19 20:17:48 +00:00
}
2021-08-03 20:40:47 +00:00
req.Res <- pathDescribeRes{Err: pathErrNoOnePublishing{PathName: pa.name}}
2020-10-19 20:17:48 +00:00
}
2021-08-11 10:45:53 +00:00
func (pa *path) handlePublisherRemove(req pathPublisherRemoveReq) {
2021-08-01 14:37:37 +00:00
if pa.source == req.Author {
2021-08-03 20:40:47 +00:00
pa.doPublisherRemove()
2021-08-01 14:37:37 +00:00
}
close(req.Res)
}
2021-08-11 10:45:53 +00:00
func (pa *path) handlePublisherAnnounce(req pathPublisherAnnounceReq) {
2021-08-01 14:37:37 +00:00
if pa.source != nil {
if pa.hasStaticSource() {
req.Res <- pathPublisherAnnounceRes{Err: fmt.Errorf("path '%s' is assigned to a static source", pa.name)}
return
}
if pa.conf.DisablePublisherOverride {
req.Res <- pathPublisherAnnounceRes{Err: fmt.Errorf("another publisher is already publishing to path '%s'", pa.name)}
return
}
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "closing existing publisher")
pa.source.(publisher).close()
2021-08-03 20:40:47 +00:00
pa.doPublisherRemove()
2021-08-01 14:37:37 +00:00
}
pa.source = req.Author
2021-08-03 20:40:47 +00:00
2021-08-01 14:37:37 +00:00
req.Res <- pathPublisherAnnounceRes{Path: pa}
}
2021-08-11 10:45:53 +00:00
func (pa *path) handlePublisherRecord(req pathPublisherRecordReq) {
2021-08-01 14:37:37 +00:00
if pa.source != req.Author {
req.Res <- pathPublisherRecordRes{Err: fmt.Errorf("publisher is not assigned to this path anymore")}
return
}
2021-10-27 19:01:00 +00:00
req.Author.onPublisherAccepted(len(req.Tracks))
2021-08-01 14:37:37 +00:00
pa.sourceSetReady(req.Tracks)
2021-08-01 14:37:37 +00:00
if pa.conf.RunOnPublish != "" {
2021-10-27 19:01:00 +00:00
pa.log(logger.Info, "runOnPublish command started")
2021-08-01 14:37:37 +00:00
_, port, _ := net.SplitHostPort(pa.rtspAddress)
pa.onPublishCmd = externalcmd.New(pa.conf.RunOnPublish, pa.conf.RunOnPublishRestart, externalcmd.Environment{
Path: pa.name,
Port: port,
})
}
req.Res <- pathPublisherRecordRes{Stream: pa.stream}
2021-08-01 14:37:37 +00:00
}
2021-08-11 10:45:53 +00:00
func (pa *path) handlePublisherPause(req pathPublisherPauseReq) {
2021-08-03 20:40:47 +00:00
if req.Author == pa.source && pa.sourceReady {
if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial {
pa.onDemandCloseSource()
} else {
pa.sourceSetNotReady()
}
2021-08-01 14:37:37 +00:00
}
close(req.Res)
}
2021-08-11 10:45:53 +00:00
func (pa *path) handleReaderRemove(req pathReaderRemoveReq) {
2021-08-01 14:37:37 +00:00
if _, ok := pa.readers[req.Author]; ok {
2021-08-03 20:40:47 +00:00
pa.doReaderRemove(req.Author)
2021-08-01 14:37:37 +00:00
}
close(req.Res)
2021-08-03 20:40:47 +00:00
if pa.isOnDemand() &&
len(pa.readers) == 0 &&
pa.onDemandState == pathOnDemandStateReady {
pa.onDemandScheduleClose()
}
2021-08-01 14:37:37 +00:00
}
2021-08-11 10:45:53 +00:00
func (pa *path) handleReaderSetupPlay(req pathReaderSetupPlayReq) {
2021-08-03 20:40:47 +00:00
if pa.sourceReady {
2021-08-11 10:45:53 +00:00
pa.handleReaderSetupPlayPost(req)
2021-03-10 14:06:45 +00:00
return
2021-08-03 20:40:47 +00:00
}
2021-03-10 14:06:45 +00:00
2021-08-03 20:40:47 +00:00
if pa.isOnDemand() {
if pa.onDemandState == pathOnDemandStateInitial {
pa.onDemandStartSource()
}
2021-03-10 14:06:45 +00:00
pa.setupPlayRequests = append(pa.setupPlayRequests, req)
return
2020-10-19 20:17:48 +00:00
}
2021-08-03 20:40:47 +00:00
req.Res <- pathReaderSetupPlayRes{Err: pathErrNoOnePublishing{PathName: pa.name}}
2021-03-10 14:06:45 +00:00
}
2020-10-19 20:17:48 +00:00
2021-08-11 10:45:53 +00:00
func (pa *path) handleReaderSetupPlayPost(req pathReaderSetupPlayReq) {
2021-08-03 20:40:47 +00:00
pa.readers[req.Author] = pathReaderStatePrePlay
2021-08-03 20:40:47 +00:00
if pa.isOnDemand() && pa.onDemandState == pathOnDemandStateClosing {
pa.onDemandState = pathOnDemandStateReady
pa.onDemandCloseTimer.Stop()
pa.onDemandCloseTimer = newEmptyTimer()
}
req.Res <- pathReaderSetupPlayRes{
Path: pa,
Stream: pa.stream,
}
2020-10-19 20:17:48 +00:00
}
2021-08-11 10:45:53 +00:00
func (pa *path) handleReaderPlay(req pathReaderPlayReq) {
pa.readers[req.Author] = pathReaderStatePlay
pa.stream.readerAdd(req.Author)
2021-10-27 19:01:00 +00:00
req.Author.onReaderAccepted()
close(req.Res)
2020-10-19 20:17:48 +00:00
}
2021-08-11 10:45:53 +00:00
func (pa *path) handleReaderPause(req pathReaderPauseReq) {
if state, ok := pa.readers[req.Author]; ok && state == pathReaderStatePlay {
pa.readers[req.Author] = pathReaderStatePrePlay
pa.stream.readerRemove(req.Author)
}
close(req.Res)
}
2020-11-07 21:47:10 +00:00
2021-11-05 16:14:31 +00:00
func (pa *path) handleAPIPathsList(req apiPathsListSubReq) {
req.Data.Items[pa.name] = apiPathsListItem{
2021-08-11 10:45:53 +00:00
ConfName: pa.confName,
Conf: pa.conf,
Source: func() interface{} {
if pa.source == nil {
return nil
}
2021-10-27 19:01:00 +00:00
return pa.source.onSourceAPIDescribe()
2021-08-11 10:45:53 +00:00
}(),
SourceReady: pa.sourceReady,
Readers: func() []interface{} {
ret := []interface{}{}
for r := range pa.readers {
2021-10-27 19:01:00 +00:00
ret = append(ret, r.onReaderAPIDescribe())
2021-08-11 10:45:53 +00:00
}
return ret
}(),
}
close(req.Res)
2021-08-11 10:45:53 +00:00
}
2021-10-27 19:01:00 +00:00
// onSourceStaticSetReady is called by a sourceStatic.
func (pa *path) onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes {
req.Res = make(chan pathSourceStaticSetReadyRes)
2021-05-10 19:32:59 +00:00
select {
case pa.sourceStaticSetReady <- req:
return <-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
return pathSourceStaticSetReadyRes{Err: fmt.Errorf("terminated")}
2021-05-10 19:32:59 +00:00
}
}
// OnSourceStaticSetNotReady is called by a sourceStatic.
func (pa *path) OnSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) {
req.Res = make(chan struct{})
2021-05-10 19:32:59 +00:00
select {
case pa.sourceStaticSetNotReady <- req:
<-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
}
2020-11-05 11:30:25 +00:00
}
2021-10-27 19:01:00 +00:00
// onDescribe is called by a reader or publisher through pathManager.
func (pa *path) onDescribe(req pathDescribeReq) pathDescribeRes {
2021-05-10 19:32:59 +00:00
select {
case pa.describe <- req:
2021-08-01 15:22:28 +00:00
return <-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
2021-08-01 15:22:28 +00:00
return pathDescribeRes{Err: fmt.Errorf("terminated")}
2021-05-10 19:32:59 +00:00
}
2020-10-19 20:17:48 +00:00
}
2021-10-27 19:01:00 +00:00
// onPublisherRemove is called by a publisher.
func (pa *path) onPublisherRemove(req pathPublisherRemoveReq) {
req.Res = make(chan struct{})
2021-05-10 19:32:59 +00:00
select {
case pa.publisherRemove <- req:
<-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
}
2020-10-19 20:17:48 +00:00
}
2021-10-27 19:01:00 +00:00
// onPublisherAnnounce is called by a publisher through pathManager.
func (pa *path) onPublisherAnnounce(req pathPublisherAnnounceReq) pathPublisherAnnounceRes {
2021-05-10 19:32:59 +00:00
select {
case pa.publisherAnnounce <- req:
2021-08-01 15:22:28 +00:00
return <-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
2021-08-01 15:22:28 +00:00
return pathPublisherAnnounceRes{Err: fmt.Errorf("terminated")}
2021-05-10 19:32:59 +00:00
}
2020-10-19 20:17:48 +00:00
}
2021-10-27 19:01:00 +00:00
// onPublisherRecord is called by a publisher.
func (pa *path) onPublisherRecord(req pathPublisherRecordReq) pathPublisherRecordRes {
req.Res = make(chan pathPublisherRecordRes)
select {
case pa.publisherRecord <- req:
return <-req.Res
case <-pa.ctx.Done():
return pathPublisherRecordRes{Err: fmt.Errorf("terminated")}
}
}
2021-10-27 19:01:00 +00:00
// onPublisherPause is called by a publisher.
func (pa *path) onPublisherPause(req pathPublisherPauseReq) {
req.Res = make(chan struct{})
2021-05-10 19:32:59 +00:00
select {
case pa.publisherPause <- req:
<-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
}
2020-10-19 20:17:48 +00:00
}
2021-10-27 19:01:00 +00:00
// onReaderRemove is called by a reader.
func (pa *path) onReaderRemove(req pathReaderRemoveReq) {
req.Res = make(chan struct{})
2021-05-10 19:32:59 +00:00
select {
case pa.readerRemove <- req:
<-req.Res
case <-pa.ctx.Done():
}
}
2021-10-27 19:01:00 +00:00
// onReaderSetupPlay is called by a reader through pathManager.
func (pa *path) onReaderSetupPlay(req pathReaderSetupPlayReq) pathReaderSetupPlayRes {
select {
case pa.readerSetupPlay <- req:
2021-08-01 15:22:28 +00:00
return <-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
2021-08-01 15:22:28 +00:00
return pathReaderSetupPlayRes{Err: fmt.Errorf("terminated")}
2021-05-10 19:32:59 +00:00
}
2020-10-19 20:17:48 +00:00
}
2021-10-27 19:01:00 +00:00
// onReaderPlay is called by a reader.
func (pa *path) onReaderPlay(req pathReaderPlayReq) {
req.Res = make(chan struct{})
2021-05-10 19:32:59 +00:00
select {
case pa.readerPlay <- req:
<-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
}
}
2021-10-27 19:01:00 +00:00
// onReaderPause is called by a reader.
func (pa *path) onReaderPause(req pathReaderPauseReq) {
req.Res = make(chan struct{})
2021-05-10 19:32:59 +00:00
select {
case pa.readerPause <- req:
<-req.Res
2021-05-10 19:32:59 +00:00
case <-pa.ctx.Done():
}
2020-11-07 21:47:10 +00:00
}
2021-03-23 19:37:43 +00:00
2021-10-27 19:01:00 +00:00
// onAPIPathsList is called by api.
2021-11-05 16:14:31 +00:00
func (pa *path) onAPIPathsList(req apiPathsListSubReq) {
req.Res = make(chan struct{})
select {
case pa.apiPathsList <- req:
<-req.Res
case <-pa.ctx.Done():
}
}