Browse Source

un-capitalize private fields

pull/764/head
aler9 5 years ago
parent
commit
8ac665be87
  1. 32
      internal/core/api.go
  2. 4
      internal/core/core_test.go
  3. 102
      internal/core/hls_muxer.go
  4. 46
      internal/core/hls_server.go
  5. 14
      internal/core/hls_source.go
  6. 20
      internal/core/metrics.go
  7. 270
      internal/core/path.go
  8. 90
      internal/core/path_manager.go
  9. 58
      internal/core/rtmp_conn.go
  10. 32
      internal/core/rtmp_server.go
  11. 16
      internal/core/rtmp_source.go
  12. 44
      internal/core/rtsp_conn.go
  13. 18
      internal/core/rtsp_server.go
  14. 68
      internal/core/rtsp_session.go
  15. 28
      internal/core/rtsp_source.go

32
internal/core/api.go

@ -427,29 +427,29 @@ func (a *api) onConfigPathsDelete(ctx *gin.Context) {
func (a *api) onPathsList(ctx *gin.Context) { func (a *api) onPathsList(ctx *gin.Context) {
res := a.pathManager.onAPIPathsList(pathAPIPathsListReq{}) res := a.pathManager.onAPIPathsList(pathAPIPathsListReq{})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusInternalServerError) ctx.AbortWithStatus(http.StatusInternalServerError)
return return
} }
ctx.JSON(http.StatusOK, res.Data) ctx.JSON(http.StatusOK, res.data)
} }
func (a *api) onRTSPSessionsList(ctx *gin.Context) { func (a *api) onRTSPSessionsList(ctx *gin.Context) {
res := a.rtspServer.onAPISessionsList(rtspServerAPISessionsListReq{}) res := a.rtspServer.onAPISessionsList(rtspServerAPISessionsListReq{})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusInternalServerError) ctx.AbortWithStatus(http.StatusInternalServerError)
return return
} }
ctx.JSON(http.StatusOK, res.Data) ctx.JSON(http.StatusOK, res.data)
} }
func (a *api) onRTSPSessionsKick(ctx *gin.Context) { func (a *api) onRTSPSessionsKick(ctx *gin.Context) {
id := ctx.Param("id") id := ctx.Param("id")
res := a.rtspServer.onAPISessionsKick(rtspServerAPISessionsKickReq{ID: id}) res := a.rtspServer.onAPISessionsKick(rtspServerAPISessionsKickReq{id: id})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusNotFound) ctx.AbortWithStatus(http.StatusNotFound)
return return
} }
@ -459,19 +459,19 @@ func (a *api) onRTSPSessionsKick(ctx *gin.Context) {
func (a *api) onRTSPSSessionsList(ctx *gin.Context) { func (a *api) onRTSPSSessionsList(ctx *gin.Context) {
res := a.rtspsServer.onAPISessionsList(rtspServerAPISessionsListReq{}) res := a.rtspsServer.onAPISessionsList(rtspServerAPISessionsListReq{})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusInternalServerError) ctx.AbortWithStatus(http.StatusInternalServerError)
return return
} }
ctx.JSON(http.StatusOK, res.Data) ctx.JSON(http.StatusOK, res.data)
} }
func (a *api) onRTSPSSessionsKick(ctx *gin.Context) { func (a *api) onRTSPSSessionsKick(ctx *gin.Context) {
id := ctx.Param("id") id := ctx.Param("id")
res := a.rtspsServer.onAPISessionsKick(rtspServerAPISessionsKickReq{ID: id}) res := a.rtspsServer.onAPISessionsKick(rtspServerAPISessionsKickReq{id: id})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusNotFound) ctx.AbortWithStatus(http.StatusNotFound)
return return
} }
@ -481,19 +481,19 @@ func (a *api) onRTSPSSessionsKick(ctx *gin.Context) {
func (a *api) onRTMPConnsList(ctx *gin.Context) { func (a *api) onRTMPConnsList(ctx *gin.Context) {
res := a.rtmpServer.onAPIConnsList(rtmpServerAPIConnsListReq{}) res := a.rtmpServer.onAPIConnsList(rtmpServerAPIConnsListReq{})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusInternalServerError) ctx.AbortWithStatus(http.StatusInternalServerError)
return return
} }
ctx.JSON(http.StatusOK, res.Data) ctx.JSON(http.StatusOK, res.data)
} }
func (a *api) onRTMPConnsKick(ctx *gin.Context) { func (a *api) onRTMPConnsKick(ctx *gin.Context) {
id := ctx.Param("id") id := ctx.Param("id")
res := a.rtmpServer.onAPIConnsKick(rtmpServerAPIConnsKickReq{ID: id}) res := a.rtmpServer.onAPIConnsKick(rtmpServerAPIConnsKickReq{id: id})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusNotFound) ctx.AbortWithStatus(http.StatusNotFound)
return return
} }
@ -503,12 +503,12 @@ func (a *api) onRTMPConnsKick(ctx *gin.Context) {
func (a *api) onHLSMuxersList(ctx *gin.Context) { func (a *api) onHLSMuxersList(ctx *gin.Context) {
res := a.hlsServer.onAPIHLSMuxersList(hlsServerAPIMuxersListReq{}) res := a.hlsServer.onAPIHLSMuxersList(hlsServerAPIMuxersListReq{})
if res.Err != nil { if res.err != nil {
ctx.AbortWithStatus(http.StatusInternalServerError) ctx.AbortWithStatus(http.StatusInternalServerError)
return return
} }
ctx.JSON(http.StatusOK, res.Data) ctx.JSON(http.StatusOK, res.data)
} }
// onConfReload is called by core. // onConfReload is called by core.

4
internal/core/core_test.go

@ -191,9 +191,9 @@ func TestCorePathAutoDeletion(t *testing.T) {
}() }()
res := p.pathManager.onAPIPathsList(pathAPIPathsListReq{}) res := p.pathManager.onAPIPathsList(pathAPIPathsListReq{})
require.NoError(t, res.Err) require.NoError(t, res.err)
require.Equal(t, 0, len(res.Data.Items)) require.Equal(t, 0, len(res.data.Items))
}) })
} }
} }

102
internal/core/hls_muxer.go

@ -95,16 +95,16 @@ window.addEventListener('DOMContentLoaded', create);
` `
type hlsMuxerResponse struct { type hlsMuxerResponse struct {
Status int status int
Header map[string]string header map[string]string
Body io.Reader body io.Reader
} }
type hlsMuxerRequest struct { type hlsMuxerRequest struct {
Dir string dir string
File string file string
Req *http.Request req *http.Request
Res chan hlsMuxerResponse res chan hlsMuxerResponse
} }
type hlsMuxerTrackIDPayloadPair struct { type hlsMuxerTrackIDPayloadPair struct {
@ -224,21 +224,21 @@ func (m *hlsMuxer) run() {
case req := <-m.request: case req := <-m.request:
if isReady { if isReady {
req.Res <- m.handleRequest(req) req.res <- m.handleRequest(req)
} else { } else {
m.requests = append(m.requests, req) m.requests = append(m.requests, req)
} }
case req := <-m.hlsServerAPIMuxersList: case req := <-m.hlsServerAPIMuxersList:
req.Data.Items[m.name] = hlsServerAPIMuxersListItem{ req.data.Items[m.name] = hlsServerAPIMuxersListItem{
LastRequest: time.Unix(atomic.LoadInt64(m.lastRequestTime), 0).String(), LastRequest: time.Unix(atomic.LoadInt64(m.lastRequestTime), 0).String(),
} }
close(req.Res) close(req.res)
case <-innerReady: case <-innerReady:
isReady = true isReady = true
for _, req := range m.requests { for _, req := range m.requests {
req.Res <- m.handleRequest(req) req.res <- m.handleRequest(req)
} }
m.requests = nil m.requests = nil
@ -252,7 +252,7 @@ func (m *hlsMuxer) run() {
m.ctxCancel() m.ctxCancel()
for _, req := range m.requests { for _, req := range m.requests {
req.Res <- hlsMuxerResponse{Status: http.StatusNotFound} req.res <- hlsMuxerResponse{status: http.StatusNotFound}
} }
m.parent.onMuxerClose(m) m.parent.onMuxerClose(m)
@ -262,18 +262,18 @@ func (m *hlsMuxer) run() {
func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{}) error { func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{}) error {
res := m.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{ res := m.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{
Author: m, author: m,
PathName: m.pathName, pathName: m.pathName,
Authenticate: nil, authenticate: nil,
}) })
if res.Err != nil { if res.err != nil {
return res.Err return res.err
} }
m.path = res.Path m.path = res.path
defer func() { defer func() {
m.path.onReaderRemove(pathReaderRemoveReq{Author: m}) m.path.onReaderRemove(pathReaderRemoveReq{author: m})
}() }()
var videoTrack *gortsplib.Track var videoTrack *gortsplib.Track
@ -283,7 +283,7 @@ func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
audioTrackID := -1 audioTrackID := -1
var aacDecoder *rtpaac.Decoder var aacDecoder *rtpaac.Decoder
for i, t := range res.Stream.tracks() { for i, t := range res.stream.tracks() {
if t.IsH264() { if t.IsH264() {
if videoTrack != nil { if videoTrack != nil {
return fmt.Errorf("can't read track %d with HLS: too many tracks", i+1) return fmt.Errorf("can't read track %d with HLS: too many tracks", i+1)
@ -330,7 +330,7 @@ func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
m.ringBuffer = ringbuffer.New(uint64(m.readBufferCount)) m.ringBuffer = ringbuffer.New(uint64(m.readBufferCount))
m.path.onReaderPlay(pathReaderPlayReq{Author: m}) m.path.onReaderPlay(pathReaderPlayReq{author: m})
writerDone := make(chan error) writerDone := make(chan error)
go func() { go func() {
@ -415,67 +415,67 @@ func (m *hlsMuxer) runInner(innerCtx context.Context, innerReady chan struct{})
func (m *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse { func (m *hlsMuxer) handleRequest(req hlsMuxerRequest) hlsMuxerResponse {
atomic.StoreInt64(m.lastRequestTime, time.Now().Unix()) atomic.StoreInt64(m.lastRequestTime, time.Now().Unix())
err := m.authenticate(req.Req) err := m.authenticate(req.req)
if err != nil { if err != nil {
if terr, ok := err.(pathErrAuthCritical); ok { if terr, ok := err.(pathErrAuthCritical); ok {
m.log(logger.Info, "authentication error: %s", terr.Message) m.log(logger.Info, "authentication error: %s", terr.message)
return hlsMuxerResponse{ return hlsMuxerResponse{
Status: http.StatusUnauthorized, status: http.StatusUnauthorized,
} }
} }
return hlsMuxerResponse{ return hlsMuxerResponse{
Status: http.StatusUnauthorized, status: http.StatusUnauthorized,
Header: map[string]string{ header: map[string]string{
"WWW-Authenticate": `Basic realm="rtsp-simple-server"`, "WWW-Authenticate": `Basic realm="rtsp-simple-server"`,
}, },
} }
} }
switch { switch {
case req.File == "index.m3u8": case req.file == "index.m3u8":
return hlsMuxerResponse{ return hlsMuxerResponse{
Status: http.StatusOK, status: http.StatusOK,
Header: map[string]string{ header: map[string]string{
"Content-Type": `application/x-mpegURL`, "Content-Type": `application/x-mpegURL`,
}, },
Body: m.muxer.PrimaryPlaylist(), body: m.muxer.PrimaryPlaylist(),
} }
case req.File == "stream.m3u8": case req.file == "stream.m3u8":
return hlsMuxerResponse{ return hlsMuxerResponse{
Status: http.StatusOK, status: http.StatusOK,
Header: map[string]string{ header: map[string]string{
"Content-Type": `application/x-mpegURL`, "Content-Type": `application/x-mpegURL`,
}, },
Body: m.muxer.StreamPlaylist(), body: m.muxer.StreamPlaylist(),
} }
case strings.HasSuffix(req.File, ".ts"): case strings.HasSuffix(req.file, ".ts"):
r := m.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}
} }
return hlsMuxerResponse{ return hlsMuxerResponse{
Status: http.StatusOK, status: http.StatusOK,
Header: map[string]string{ header: map[string]string{
"Content-Type": `video/MP2T`, "Content-Type": `video/MP2T`,
}, },
Body: r, body: r,
} }
case req.File == "": case req.file == "":
return hlsMuxerResponse{ return hlsMuxerResponse{
Status: http.StatusOK, status: http.StatusOK,
Header: map[string]string{ header: map[string]string{
"Content-Type": `text/html`, "Content-Type": `text/html`,
}, },
Body: bytes.NewReader([]byte(index)), body: bytes.NewReader([]byte(index)),
} }
default: default:
return hlsMuxerResponse{Status: http.StatusNotFound} return hlsMuxerResponse{status: http.StatusNotFound}
} }
} }
@ -499,7 +499,7 @@ func (m *hlsMuxer) authenticate(req *http.Request) error {
"read") "read")
if err != nil { if err != nil {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: fmt.Sprintf("external authentication failed: %s", err), message: fmt.Sprintf("external authentication failed: %s", err),
} }
} }
} }
@ -510,7 +510,7 @@ func (m *hlsMuxer) authenticate(req *http.Request) error {
if !ipEqualOrInRange(ip, pathIPs) { if !ipEqualOrInRange(ip, pathIPs) {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: fmt.Sprintf("IP '%s' not allowed", ip), message: fmt.Sprintf("IP '%s' not allowed", ip),
} }
} }
} }
@ -523,7 +523,7 @@ func (m *hlsMuxer) authenticate(req *http.Request) error {
if user != string(pathUser) || pass != string(pathPass) { if user != string(pathUser) || pass != string(pathPass) {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: "invalid credentials", message: "invalid credentials",
} }
} }
} }
@ -536,7 +536,7 @@ func (m *hlsMuxer) onRequest(req hlsMuxerRequest) {
select { select {
case m.request <- req: case m.request <- req:
case <-m.ctx.Done(): case <-m.ctx.Done():
req.Res <- hlsMuxerResponse{Status: http.StatusNotFound} req.res <- hlsMuxerResponse{status: http.StatusNotFound}
} }
} }
@ -563,10 +563,10 @@ func (m *hlsMuxer) onReaderAPIDescribe() interface{} {
// onAPIHLSMuxersList is called by api. // onAPIHLSMuxersList is called by api.
func (m *hlsMuxer) onAPIHLSMuxersList(req hlsServerAPIMuxersListSubReq) { func (m *hlsMuxer) onAPIHLSMuxersList(req hlsServerAPIMuxersListSubReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case m.hlsServerAPIMuxersList <- req: case m.hlsServerAPIMuxersList <- req:
<-req.Res <-req.res
case <-m.ctx.Done(): case <-m.ctx.Done():
} }

46
internal/core/hls_server.go

@ -26,18 +26,18 @@ type hlsServerAPIMuxersListData struct {
} }
type hlsServerAPIMuxersListRes struct { type hlsServerAPIMuxersListRes struct {
Data *hlsServerAPIMuxersListData data *hlsServerAPIMuxersListData
Muxers map[string]*hlsMuxer muxers map[string]*hlsMuxer
Err error err error
} }
type hlsServerAPIMuxersListReq struct { type hlsServerAPIMuxersListReq struct {
Res chan hlsServerAPIMuxersListRes res chan hlsServerAPIMuxersListRes
} }
type hlsServerAPIMuxersListSubReq struct { type hlsServerAPIMuxersListSubReq struct {
Data *hlsServerAPIMuxersListData data *hlsServerAPIMuxersListData
Res chan struct{} res chan struct{}
} }
type hlsServerParent interface { type hlsServerParent interface {
@ -151,7 +151,7 @@ outer:
} }
case req := <-s.request: case req := <-s.request:
r := s.findOrCreateMuxer(req.Dir) r := s.findOrCreateMuxer(req.dir)
r.onRequest(req) r.onRequest(req)
case c := <-s.muxerClose: case c := <-s.muxerClose:
@ -167,8 +167,8 @@ outer:
muxers[name] = m muxers[name] = m
} }
req.Res <- hlsServerAPIMuxersListRes{ req.res <- hlsServerAPIMuxersListRes{
Muxers: muxers, muxers: muxers,
} }
case <-s.ctx.Done(): case <-s.ctx.Done():
@ -240,23 +240,23 @@ func (s *hlsServer) onRequest(ctx *gin.Context) {
cres := make(chan hlsMuxerResponse) cres := make(chan hlsMuxerResponse)
hreq := hlsMuxerRequest{ hreq := hlsMuxerRequest{
Dir: dir, dir: dir,
File: fname, file: fname,
Req: ctx.Request, req: ctx.Request,
Res: cres, res: cres,
} }
select { select {
case s.request <- hreq: case s.request <- hreq:
res := <-cres res := <-cres
for k, v := range res.Header { for k, v := range res.header {
ctx.Writer.Header().Set(k, v) ctx.Writer.Header().Set(k, v)
} }
ctx.Writer.WriteHeader(res.Status) ctx.Writer.WriteHeader(res.status)
if res.Body != nil { if res.body != nil {
io.Copy(ctx.Writer, res.Body) io.Copy(ctx.Writer, res.body)
} }
case <-s.ctx.Done(): case <-s.ctx.Done():
@ -303,22 +303,22 @@ func (s *hlsServer) onPathSourceReady(pa *path) {
// onAPIHLSMuxersList is called by api. // onAPIHLSMuxersList is called by api.
func (s *hlsServer) onAPIHLSMuxersList(req hlsServerAPIMuxersListReq) hlsServerAPIMuxersListRes { func (s *hlsServer) onAPIHLSMuxersList(req hlsServerAPIMuxersListReq) hlsServerAPIMuxersListRes {
req.Res = make(chan hlsServerAPIMuxersListRes) req.res = make(chan hlsServerAPIMuxersListRes)
select { select {
case s.apiMuxersList <- req: case s.apiMuxersList <- req:
res := <-req.Res res := <-req.res
res.Data = &hlsServerAPIMuxersListData{ res.data = &hlsServerAPIMuxersListData{
Items: make(map[string]hlsServerAPIMuxersListItem), Items: make(map[string]hlsServerAPIMuxersListItem),
} }
for _, pa := range res.Muxers { for _, pa := range res.muxers {
pa.onAPIHLSMuxersList(hlsServerAPIMuxersListSubReq{Data: res.Data}) pa.onAPIHLSMuxersList(hlsServerAPIMuxersListSubReq{data: res.data})
} }
return res return res
case <-s.ctx.Done(): case <-s.ctx.Done():
return hlsServerAPIMuxersListRes{Err: fmt.Errorf("terminated")} return hlsServerAPIMuxersListRes{err: fmt.Errorf("terminated")}
} }
} }

14
internal/core/hls_source.go

@ -19,7 +19,7 @@ const (
type hlsSourceParent interface { type hlsSourceParent interface {
log(logger.Level, string, ...interface{}) log(logger.Level, string, ...interface{})
onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes
OnSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) onSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq)
} }
type hlsSource struct { type hlsSource struct {
@ -94,7 +94,7 @@ func (s *hlsSource) runInner() bool {
defer func() { defer func() {
if stream != nil { if stream != nil {
s.parent.OnSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{Source: s}) s.parent.onSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{source: s})
rtcpSenders.Close() rtcpSenders.Close()
} }
}() }()
@ -113,16 +113,16 @@ func (s *hlsSource) runInner() bool {
} }
res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{ res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{
Source: s, source: s,
Tracks: tracks, tracks: tracks,
}) })
if res.Err != nil { if res.err != nil {
return res.Err return res.err
} }
s.Log(logger.Info, "ready") s.Log(logger.Info, "ready")
stream = res.Stream stream = res.stream
rtcpSenders = rtcpsenderset.New(tracks, stream.onPacketRTCP) rtcpSenders = rtcpsenderset.New(tracks, stream.onPacketRTCP)
return nil return nil

20
internal/core/metrics.go

@ -96,8 +96,8 @@ func (m *metrics) onMetrics(ctx *gin.Context) {
out := "" out := ""
res := m.pathManager.onAPIPathsList(pathAPIPathsListReq{}) res := m.pathManager.onAPIPathsList(pathAPIPathsListReq{})
if res.Err == nil { if res.err == nil {
for name, p := range res.Data.Items { for name, p := range res.data.Items {
if p.SourceReady { if p.SourceReady {
out += metric("paths{name=\""+name+"\",state=\"ready\"}", 1) out += metric("paths{name=\""+name+"\",state=\"ready\"}", 1)
} else { } else {
@ -108,12 +108,12 @@ func (m *metrics) onMetrics(ctx *gin.Context) {
if !interfaceIsEmpty(m.rtspServer) { if !interfaceIsEmpty(m.rtspServer) {
res := m.rtspServer.onAPISessionsList(rtspServerAPISessionsListReq{}) res := m.rtspServer.onAPISessionsList(rtspServerAPISessionsListReq{})
if res.Err == nil { if res.err == nil {
idleCount := int64(0) idleCount := int64(0)
readCount := int64(0) readCount := int64(0)
publishCount := int64(0) publishCount := int64(0)
for _, i := range res.Data.Items { for _, i := range res.data.Items {
switch i.State { switch i.State {
case "idle": case "idle":
idleCount++ idleCount++
@ -135,12 +135,12 @@ func (m *metrics) onMetrics(ctx *gin.Context) {
if !interfaceIsEmpty(m.rtspsServer) { if !interfaceIsEmpty(m.rtspsServer) {
res := m.rtspsServer.onAPISessionsList(rtspServerAPISessionsListReq{}) res := m.rtspsServer.onAPISessionsList(rtspServerAPISessionsListReq{})
if res.Err == nil { if res.err == nil {
idleCount := int64(0) idleCount := int64(0)
readCount := int64(0) readCount := int64(0)
publishCount := int64(0) publishCount := int64(0)
for _, i := range res.Data.Items { for _, i := range res.data.Items {
switch i.State { switch i.State {
case "idle": case "idle":
idleCount++ idleCount++
@ -162,12 +162,12 @@ func (m *metrics) onMetrics(ctx *gin.Context) {
if !interfaceIsEmpty(m.rtmpServer) { if !interfaceIsEmpty(m.rtmpServer) {
res := m.rtmpServer.onAPIConnsList(rtmpServerAPIConnsListReq{}) res := m.rtmpServer.onAPIConnsList(rtmpServerAPIConnsListReq{})
if res.Err == nil { if res.err == nil {
idleCount := int64(0) idleCount := int64(0)
readCount := int64(0) readCount := int64(0)
publishCount := int64(0) publishCount := int64(0)
for _, i := range res.Data.Items { for _, i := range res.data.Items {
switch i.State { switch i.State {
case "idle": case "idle":
idleCount++ idleCount++
@ -189,8 +189,8 @@ func (m *metrics) onMetrics(ctx *gin.Context) {
if !interfaceIsEmpty(m.hlsServer) { if !interfaceIsEmpty(m.hlsServer) {
res := m.hlsServer.onAPIHLSMuxersList(hlsServerAPIMuxersListReq{}) res := m.hlsServer.onAPIHLSMuxersList(hlsServerAPIMuxersListReq{})
if res.Err == nil { if res.err == nil {
for name := range res.Data.Items { for name := range res.data.Items {
out += metric("hls_muxers{name=\""+name+"\"}", 1) out += metric("hls_muxers{name=\""+name+"\"}", 1)
} }
} }

270
internal/core/path.go

@ -30,17 +30,17 @@ type authenticateFunc func(
) error ) error
type pathErrNoOnePublishing struct { type pathErrNoOnePublishing struct {
PathName string pathName string
} }
// Error implements the error interface. // Error implements the error interface.
func (e pathErrNoOnePublishing) Error() string { func (e pathErrNoOnePublishing) Error() string {
return fmt.Sprintf("no one is publishing to path '%s'", e.PathName) return fmt.Sprintf("no one is publishing to path '%s'", e.pathName)
} }
type pathErrAuthNotCritical struct { type pathErrAuthNotCritical struct {
Message string message string
Response *base.Response response *base.Response
} }
// Error implements the error interface. // Error implements the error interface.
@ -49,8 +49,8 @@ func (pathErrAuthNotCritical) Error() string {
} }
type pathErrAuthCritical struct { type pathErrAuthCritical struct {
Message string message string
Response *base.Response response *base.Response
} }
// Error implements the error interface. // Error implements the error interface.
@ -94,94 +94,94 @@ const (
) )
type pathSourceStaticSetReadyRes struct { type pathSourceStaticSetReadyRes struct {
Stream *stream stream *stream
Err error err error
} }
type pathSourceStaticSetReadyReq struct { type pathSourceStaticSetReadyReq struct {
Source sourceStatic source sourceStatic
Tracks gortsplib.Tracks tracks gortsplib.Tracks
Res chan pathSourceStaticSetReadyRes res chan pathSourceStaticSetReadyRes
} }
type pathSourceStaticSetNotReadyReq struct { type pathSourceStaticSetNotReadyReq struct {
Source sourceStatic source sourceStatic
Res chan struct{} res chan struct{}
} }
type pathReaderRemoveReq struct { type pathReaderRemoveReq struct {
Author reader author reader
Res chan struct{} res chan struct{}
} }
type pathPublisherRemoveReq struct { type pathPublisherRemoveReq struct {
Author publisher author publisher
Res chan struct{} res chan struct{}
} }
type pathDescribeRes struct { type pathDescribeRes struct {
Path *path path *path
Stream *stream stream *stream
Redirect string redirect string
Err error err error
} }
type pathDescribeReq struct { type pathDescribeReq struct {
PathName string pathName string
URL *base.URL url *base.URL
Authenticate authenticateFunc authenticate authenticateFunc
Res chan pathDescribeRes res chan pathDescribeRes
} }
type pathReaderSetupPlayRes struct { type pathReaderSetupPlayRes struct {
Path *path path *path
Stream *stream stream *stream
Err error err error
} }
type pathReaderSetupPlayReq struct { type pathReaderSetupPlayReq struct {
Author reader author reader
PathName string pathName string
Authenticate authenticateFunc authenticate authenticateFunc
Res chan pathReaderSetupPlayRes res chan pathReaderSetupPlayRes
} }
type pathPublisherAnnounceRes struct { type pathPublisherAnnounceRes struct {
Path *path path *path
Err error err error
} }
type pathPublisherAnnounceReq struct { type pathPublisherAnnounceReq struct {
Author publisher author publisher
PathName string pathName string
Authenticate authenticateFunc authenticate authenticateFunc
Res chan pathPublisherAnnounceRes res chan pathPublisherAnnounceRes
} }
type pathReaderPlayReq struct { type pathReaderPlayReq struct {
Author reader author reader
Res chan struct{} res chan struct{}
} }
type pathPublisherRecordRes struct { type pathPublisherRecordRes struct {
Stream *stream stream *stream
Err error err error
} }
type pathPublisherRecordReq struct { type pathPublisherRecordReq struct {
Author publisher author publisher
Tracks gortsplib.Tracks tracks gortsplib.Tracks
Res chan pathPublisherRecordRes res chan pathPublisherRecordRes
} }
type pathReaderPauseReq struct { type pathReaderPauseReq struct {
Author reader author reader
Res chan struct{} res chan struct{}
} }
type pathPublisherPauseReq struct { type pathPublisherPauseReq struct {
Author publisher author publisher
Res chan struct{} res chan struct{}
} }
type pathAPIPathsListItem struct { type pathAPIPathsListItem struct {
@ -197,18 +197,18 @@ type pathAPIPathsListData struct {
} }
type pathAPIPathsListRes struct { type pathAPIPathsListRes struct {
Data *pathAPIPathsListData data *pathAPIPathsListData
Paths map[string]*path paths map[string]*path
Err error err error
} }
type pathAPIPathsListReq struct { type pathAPIPathsListReq struct {
Res chan pathAPIPathsListRes res chan pathAPIPathsListRes
} }
type pathAPIPathsListSubReq struct { type pathAPIPathsListSubReq struct {
Data *pathAPIPathsListData data *pathAPIPathsListData
Res chan struct{} res chan struct{}
} }
type path struct { type path struct {
@ -362,12 +362,12 @@ func (pa *path) run() {
select { select {
case <-pa.onDemandReadyTimer.C: case <-pa.onDemandReadyTimer.C:
for _, req := range pa.describeRequests { for _, req := range pa.describeRequests {
req.Res <- pathDescribeRes{Err: fmt.Errorf("source of path '%s' has timed out", pa.name)} req.res <- pathDescribeRes{err: fmt.Errorf("source of path '%s' has timed out", pa.name)}
} }
pa.describeRequests = nil pa.describeRequests = nil
for _, req := range pa.setupPlayRequests { for _, req := range pa.setupPlayRequests {
req.Res <- pathReaderSetupPlayRes{Err: fmt.Errorf("source of path '%s' has timed out", pa.name)} req.res <- pathReaderSetupPlayRes{err: fmt.Errorf("source of path '%s' has timed out", pa.name)}
} }
pa.setupPlayRequests = nil pa.setupPlayRequests = nil
@ -385,22 +385,22 @@ func (pa *path) run() {
} }
case req := <-pa.sourceStaticSetReady: case req := <-pa.sourceStaticSetReady:
if req.Source == pa.source { if req.source == pa.source {
pa.sourceSetReady(req.Tracks) pa.sourceSetReady(req.tracks)
req.Res <- pathSourceStaticSetReadyRes{Stream: pa.stream} req.res <- pathSourceStaticSetReadyRes{stream: pa.stream}
} else { } else {
req.Res <- pathSourceStaticSetReadyRes{Err: fmt.Errorf("terminated")} req.res <- pathSourceStaticSetReadyRes{err: fmt.Errorf("terminated")}
} }
case req := <-pa.sourceStaticSetNotReady: case req := <-pa.sourceStaticSetNotReady:
if req.Source == pa.source { if req.source == pa.source {
if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial { if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial {
pa.onDemandCloseSource() pa.onDemandCloseSource()
} else { } else {
pa.sourceSetNotReady() pa.sourceSetNotReady()
} }
} }
close(req.Res) close(req.res)
if pa.shouldClose() { if pa.shouldClose() {
return fmt.Errorf("not in use") return fmt.Errorf("not in use")
@ -469,11 +469,11 @@ func (pa *path) run() {
} }
for _, req := range pa.describeRequests { for _, req := range pa.describeRequests {
req.Res <- pathDescribeRes{Err: fmt.Errorf("terminated")} req.res <- pathDescribeRes{err: fmt.Errorf("terminated")}
} }
for _, req := range pa.setupPlayRequests { for _, req := range pa.setupPlayRequests {
req.Res <- pathReaderSetupPlayRes{Err: fmt.Errorf("terminated")} req.res <- pathReaderSetupPlayRes{err: fmt.Errorf("terminated")}
} }
for rp := range pa.readers { for rp := range pa.readers {
@ -609,8 +609,8 @@ func (pa *path) sourceSetReady(tracks gortsplib.Tracks) {
pa.onDemandReadyTimer = newEmptyTimer() pa.onDemandReadyTimer = newEmptyTimer()
for _, req := range pa.describeRequests { for _, req := range pa.describeRequests {
req.Res <- pathDescribeRes{ req.res <- pathDescribeRes{
Stream: pa.stream, stream: pa.stream,
} }
} }
pa.describeRequests = nil pa.describeRequests = nil
@ -711,15 +711,15 @@ func (pa *path) doPublisherRemove() {
func (pa *path) handleDescribe(req pathDescribeReq) { func (pa *path) handleDescribe(req pathDescribeReq) {
if _, ok := pa.source.(*sourceRedirect); ok { if _, ok := pa.source.(*sourceRedirect); ok {
req.Res <- pathDescribeRes{ req.res <- pathDescribeRes{
Redirect: pa.conf.SourceRedirect, redirect: pa.conf.SourceRedirect,
} }
return return
} }
if pa.sourceReady { if pa.sourceReady {
req.Res <- pathDescribeRes{ req.res <- pathDescribeRes{
Stream: pa.stream, stream: pa.stream,
} }
return return
} }
@ -736,38 +736,38 @@ func (pa *path) handleDescribe(req pathDescribeReq) {
fallbackURL := func() string { fallbackURL := func() string {
if strings.HasPrefix(pa.conf.Fallback, "/") { if strings.HasPrefix(pa.conf.Fallback, "/") {
ur := base.URL{ ur := base.URL{
Scheme: req.URL.Scheme, Scheme: req.url.Scheme,
User: req.URL.User, User: req.url.User,
Host: req.URL.Host, Host: req.url.Host,
Path: pa.conf.Fallback, Path: pa.conf.Fallback,
} }
return ur.String() return ur.String()
} }
return pa.conf.Fallback return pa.conf.Fallback
}() }()
req.Res <- pathDescribeRes{Redirect: fallbackURL} req.res <- pathDescribeRes{redirect: fallbackURL}
return return
} }
req.Res <- pathDescribeRes{Err: pathErrNoOnePublishing{PathName: pa.name}} req.res <- pathDescribeRes{err: pathErrNoOnePublishing{pathName: pa.name}}
} }
func (pa *path) handlePublisherRemove(req pathPublisherRemoveReq) { func (pa *path) handlePublisherRemove(req pathPublisherRemoveReq) {
if pa.source == req.Author { if pa.source == req.author {
pa.doPublisherRemove() pa.doPublisherRemove()
} }
close(req.Res) close(req.res)
} }
func (pa *path) handlePublisherAnnounce(req pathPublisherAnnounceReq) { func (pa *path) handlePublisherAnnounce(req pathPublisherAnnounceReq) {
if pa.source != nil { if pa.source != nil {
if pa.hasStaticSource() { if pa.hasStaticSource() {
req.Res <- pathPublisherAnnounceRes{Err: fmt.Errorf("path '%s' is assigned to a static source", pa.name)} req.res <- pathPublisherAnnounceRes{err: fmt.Errorf("path '%s' is assigned to a static source", pa.name)}
return return
} }
if pa.conf.DisablePublisherOverride { if pa.conf.DisablePublisherOverride {
req.Res <- pathPublisherAnnounceRes{Err: fmt.Errorf("another publisher is already publishing to path '%s'", pa.name)} req.res <- pathPublisherAnnounceRes{err: fmt.Errorf("another publisher is already publishing to path '%s'", pa.name)}
return return
} }
@ -776,20 +776,20 @@ func (pa *path) handlePublisherAnnounce(req pathPublisherAnnounceReq) {
pa.doPublisherRemove() pa.doPublisherRemove()
} }
pa.source = req.Author pa.source = req.author
req.Res <- pathPublisherAnnounceRes{Path: pa} req.res <- pathPublisherAnnounceRes{path: pa}
} }
func (pa *path) handlePublisherRecord(req pathPublisherRecordReq) { func (pa *path) handlePublisherRecord(req pathPublisherRecordReq) {
if pa.source != req.Author { if pa.source != req.author {
req.Res <- pathPublisherRecordRes{Err: fmt.Errorf("publisher is not assigned to this path anymore")} req.res <- pathPublisherRecordRes{err: fmt.Errorf("publisher is not assigned to this path anymore")}
return return
} }
req.Author.onPublisherAccepted(len(req.Tracks)) req.author.onPublisherAccepted(len(req.tracks))
pa.sourceSetReady(req.Tracks) pa.sourceSetReady(req.tracks)
if pa.conf.RunOnPublish != "" { if pa.conf.RunOnPublish != "" {
pa.log(logger.Info, "runOnPublish command started") pa.log(logger.Info, "runOnPublish command started")
@ -803,25 +803,25 @@ func (pa *path) handlePublisherRecord(req pathPublisherRecordReq) {
}) })
} }
req.Res <- pathPublisherRecordRes{Stream: pa.stream} req.res <- pathPublisherRecordRes{stream: pa.stream}
} }
func (pa *path) handlePublisherPause(req pathPublisherPauseReq) { func (pa *path) handlePublisherPause(req pathPublisherPauseReq) {
if req.Author == pa.source && pa.sourceReady { if req.author == pa.source && pa.sourceReady {
if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial { if pa.isOnDemand() && pa.onDemandState != pathOnDemandStateInitial {
pa.onDemandCloseSource() pa.onDemandCloseSource()
} else { } else {
pa.sourceSetNotReady() pa.sourceSetNotReady()
} }
} }
close(req.Res) close(req.res)
} }
func (pa *path) handleReaderRemove(req pathReaderRemoveReq) { func (pa *path) handleReaderRemove(req pathReaderRemoveReq) {
if _, ok := pa.readers[req.Author]; ok { if _, ok := pa.readers[req.author]; ok {
pa.doReaderRemove(req.Author) pa.doReaderRemove(req.author)
} }
close(req.Res) close(req.res)
if pa.isOnDemand() && if pa.isOnDemand() &&
len(pa.readers) == 0 && len(pa.readers) == 0 &&
@ -844,11 +844,11 @@ func (pa *path) handleReaderSetupPlay(req pathReaderSetupPlayReq) {
return return
} }
req.Res <- pathReaderSetupPlayRes{Err: pathErrNoOnePublishing{PathName: pa.name}} req.res <- pathReaderSetupPlayRes{err: pathErrNoOnePublishing{pathName: pa.name}}
} }
func (pa *path) handleReaderSetupPlayPost(req pathReaderSetupPlayReq) { func (pa *path) handleReaderSetupPlayPost(req pathReaderSetupPlayReq) {
pa.readers[req.Author] = pathReaderStatePrePlay pa.readers[req.author] = pathReaderStatePrePlay
if pa.isOnDemand() && pa.onDemandState == pathOnDemandStateClosing { if pa.isOnDemand() && pa.onDemandState == pathOnDemandStateClosing {
pa.onDemandState = pathOnDemandStateReady pa.onDemandState = pathOnDemandStateReady
@ -856,32 +856,32 @@ func (pa *path) handleReaderSetupPlayPost(req pathReaderSetupPlayReq) {
pa.onDemandCloseTimer = newEmptyTimer() pa.onDemandCloseTimer = newEmptyTimer()
} }
req.Res <- pathReaderSetupPlayRes{ req.res <- pathReaderSetupPlayRes{
Path: pa, path: pa,
Stream: pa.stream, stream: pa.stream,
} }
} }
func (pa *path) handleReaderPlay(req pathReaderPlayReq) { func (pa *path) handleReaderPlay(req pathReaderPlayReq) {
pa.readers[req.Author] = pathReaderStatePlay pa.readers[req.author] = pathReaderStatePlay
pa.stream.readerAdd(req.Author) pa.stream.readerAdd(req.author)
req.Author.onReaderAccepted() req.author.onReaderAccepted()
close(req.Res) close(req.res)
} }
func (pa *path) handleReaderPause(req pathReaderPauseReq) { func (pa *path) handleReaderPause(req pathReaderPauseReq) {
if state, ok := pa.readers[req.Author]; ok && state == pathReaderStatePlay { if state, ok := pa.readers[req.author]; ok && state == pathReaderStatePlay {
pa.readers[req.Author] = pathReaderStatePrePlay pa.readers[req.author] = pathReaderStatePrePlay
pa.stream.readerRemove(req.Author) pa.stream.readerRemove(req.author)
} }
close(req.Res) close(req.res)
} }
func (pa *path) handleAPIPathsList(req pathAPIPathsListSubReq) { func (pa *path) handleAPIPathsList(req pathAPIPathsListSubReq) {
req.Data.Items[pa.name] = pathAPIPathsListItem{ req.data.Items[pa.name] = pathAPIPathsListItem{
ConfName: pa.confName, ConfName: pa.confName,
Conf: pa.conf, Conf: pa.conf,
Source: func() interface{} { Source: func() interface{} {
@ -899,26 +899,26 @@ func (pa *path) handleAPIPathsList(req pathAPIPathsListSubReq) {
return ret return ret
}(), }(),
} }
close(req.Res) close(req.res)
} }
// onSourceStaticSetReady is called by a sourceStatic. // onSourceStaticSetReady is called by a sourceStatic.
func (pa *path) onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes { func (pa *path) onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes {
req.Res = make(chan pathSourceStaticSetReadyRes) req.res = make(chan pathSourceStaticSetReadyRes)
select { select {
case pa.sourceStaticSetReady <- req: case pa.sourceStaticSetReady <- req:
return <-req.Res return <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
return pathSourceStaticSetReadyRes{Err: fmt.Errorf("terminated")} return pathSourceStaticSetReadyRes{err: fmt.Errorf("terminated")}
} }
} }
// OnSourceStaticSetNotReady is called by a sourceStatic. // onSourceStaticSetNotReady is called by a sourceStatic.
func (pa *path) OnSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) { func (pa *path) onSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.sourceStaticSetNotReady <- req: case pa.sourceStaticSetNotReady <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }
} }
@ -927,18 +927,18 @@ func (pa *path) OnSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) {
func (pa *path) onDescribe(req pathDescribeReq) pathDescribeRes { func (pa *path) onDescribe(req pathDescribeReq) pathDescribeRes {
select { select {
case pa.describe <- req: case pa.describe <- req:
return <-req.Res return <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
return pathDescribeRes{Err: fmt.Errorf("terminated")} return pathDescribeRes{err: fmt.Errorf("terminated")}
} }
} }
// onPublisherRemove is called by a publisher. // onPublisherRemove is called by a publisher.
func (pa *path) onPublisherRemove(req pathPublisherRemoveReq) { func (pa *path) onPublisherRemove(req pathPublisherRemoveReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.publisherRemove <- req: case pa.publisherRemove <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }
} }
@ -947,39 +947,39 @@ func (pa *path) onPublisherRemove(req pathPublisherRemoveReq) {
func (pa *path) onPublisherAnnounce(req pathPublisherAnnounceReq) pathPublisherAnnounceRes { func (pa *path) onPublisherAnnounce(req pathPublisherAnnounceReq) pathPublisherAnnounceRes {
select { select {
case pa.publisherAnnounce <- req: case pa.publisherAnnounce <- req:
return <-req.Res return <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
return pathPublisherAnnounceRes{Err: fmt.Errorf("terminated")} return pathPublisherAnnounceRes{err: fmt.Errorf("terminated")}
} }
} }
// onPublisherRecord is called by a publisher. // onPublisherRecord is called by a publisher.
func (pa *path) onPublisherRecord(req pathPublisherRecordReq) pathPublisherRecordRes { func (pa *path) onPublisherRecord(req pathPublisherRecordReq) pathPublisherRecordRes {
req.Res = make(chan pathPublisherRecordRes) req.res = make(chan pathPublisherRecordRes)
select { select {
case pa.publisherRecord <- req: case pa.publisherRecord <- req:
return <-req.Res return <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
return pathPublisherRecordRes{Err: fmt.Errorf("terminated")} return pathPublisherRecordRes{err: fmt.Errorf("terminated")}
} }
} }
// onPublisherPause is called by a publisher. // onPublisherPause is called by a publisher.
func (pa *path) onPublisherPause(req pathPublisherPauseReq) { func (pa *path) onPublisherPause(req pathPublisherPauseReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.publisherPause <- req: case pa.publisherPause <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }
} }
// onReaderRemove is called by a reader. // onReaderRemove is called by a reader.
func (pa *path) onReaderRemove(req pathReaderRemoveReq) { func (pa *path) onReaderRemove(req pathReaderRemoveReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.readerRemove <- req: case pa.readerRemove <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }
} }
@ -988,38 +988,38 @@ func (pa *path) onReaderRemove(req pathReaderRemoveReq) {
func (pa *path) onReaderSetupPlay(req pathReaderSetupPlayReq) pathReaderSetupPlayRes { func (pa *path) onReaderSetupPlay(req pathReaderSetupPlayReq) pathReaderSetupPlayRes {
select { select {
case pa.readerSetupPlay <- req: case pa.readerSetupPlay <- req:
return <-req.Res return <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
return pathReaderSetupPlayRes{Err: fmt.Errorf("terminated")} return pathReaderSetupPlayRes{err: fmt.Errorf("terminated")}
} }
} }
// onReaderPlay is called by a reader. // onReaderPlay is called by a reader.
func (pa *path) onReaderPlay(req pathReaderPlayReq) { func (pa *path) onReaderPlay(req pathReaderPlayReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.readerPlay <- req: case pa.readerPlay <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }
} }
// onReaderPause is called by a reader. // onReaderPause is called by a reader.
func (pa *path) onReaderPause(req pathReaderPauseReq) { func (pa *path) onReaderPause(req pathReaderPauseReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.readerPause <- req: case pa.readerPause <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }
} }
// onAPIPathsList is called by api. // onAPIPathsList is called by api.
func (pa *path) onAPIPathsList(req pathAPIPathsListSubReq) { func (pa *path) onAPIPathsList(req pathAPIPathsListSubReq) {
req.Res = make(chan struct{}) req.res = make(chan struct{})
select { select {
case pa.apiPathsList <- req: case pa.apiPathsList <- req:
<-req.Res <-req.res
case <-pa.ctx.Done(): case <-pa.ctx.Done():
} }

90
internal/core/path_manager.go

@ -165,75 +165,75 @@ outer:
} }
case req := <-pm.describe: case req := <-pm.describe:
pathConfName, pathConf, pathMatches, err := pm.findPathConf(req.PathName) pathConfName, pathConf, pathMatches, err := pm.findPathConf(req.pathName)
if err != nil { if err != nil {
req.Res <- pathDescribeRes{Err: err} req.res <- pathDescribeRes{err: err}
continue continue
} }
err = req.Authenticate( err = req.authenticate(
pathConf.ReadIPs, pathConf.ReadIPs,
pathConf.ReadUser, pathConf.ReadUser,
pathConf.ReadPass) pathConf.ReadPass)
if err != nil { if err != nil {
req.Res <- pathDescribeRes{Err: err} req.res <- pathDescribeRes{err: err}
continue continue
} }
// create path if it doesn't exist // create path if it doesn't exist
if _, ok := pm.paths[req.PathName]; !ok { if _, ok := pm.paths[req.pathName]; !ok {
pm.createPath(pathConfName, pathConf, req.PathName, pathMatches) pm.createPath(pathConfName, pathConf, req.pathName, pathMatches)
} }
req.Res <- pathDescribeRes{Path: pm.paths[req.PathName]} req.res <- pathDescribeRes{path: pm.paths[req.pathName]}
case req := <-pm.readerSetupPlay: case req := <-pm.readerSetupPlay:
pathConfName, pathConf, pathMatches, err := pm.findPathConf(req.PathName) pathConfName, pathConf, pathMatches, err := pm.findPathConf(req.pathName)
if err != nil { if err != nil {
req.Res <- pathReaderSetupPlayRes{Err: err} req.res <- pathReaderSetupPlayRes{err: err}
continue continue
} }
if req.Authenticate != nil { if req.authenticate != nil {
err = req.Authenticate( err = req.authenticate(
pathConf.ReadIPs, pathConf.ReadIPs,
pathConf.ReadUser, pathConf.ReadUser,
pathConf.ReadPass) pathConf.ReadPass)
if err != nil { if err != nil {
req.Res <- pathReaderSetupPlayRes{Err: err} req.res <- pathReaderSetupPlayRes{err: err}
continue continue
} }
} }
// create path if it doesn't exist // create path if it doesn't exist
if _, ok := pm.paths[req.PathName]; !ok { if _, ok := pm.paths[req.pathName]; !ok {
pm.createPath(pathConfName, pathConf, req.PathName, pathMatches) pm.createPath(pathConfName, pathConf, req.pathName, pathMatches)
} }
req.Res <- pathReaderSetupPlayRes{Path: pm.paths[req.PathName]} req.res <- pathReaderSetupPlayRes{path: pm.paths[req.pathName]}
case req := <-pm.publisherAnnounce: case req := <-pm.publisherAnnounce:
pathConfName, pathConf, pathMatches, err := pm.findPathConf(req.PathName) pathConfName, pathConf, pathMatches, err := pm.findPathConf(req.pathName)
if err != nil { if err != nil {
req.Res <- pathPublisherAnnounceRes{Err: err} req.res <- pathPublisherAnnounceRes{err: err}
continue continue
} }
err = req.Authenticate( err = req.authenticate(
pathConf.PublishIPs, pathConf.PublishIPs,
pathConf.PublishUser, pathConf.PublishUser,
pathConf.PublishPass) pathConf.PublishPass)
if err != nil { if err != nil {
req.Res <- pathPublisherAnnounceRes{Err: err} req.res <- pathPublisherAnnounceRes{err: err}
continue continue
} }
// create path if it doesn't exist // create path if it doesn't exist
if _, ok := pm.paths[req.PathName]; !ok { if _, ok := pm.paths[req.pathName]; !ok {
pm.createPath(pathConfName, pathConf, req.PathName, pathMatches) pm.createPath(pathConfName, pathConf, req.pathName, pathMatches)
} }
req.Res <- pathPublisherAnnounceRes{Path: pm.paths[req.PathName]} req.res <- pathPublisherAnnounceRes{path: pm.paths[req.pathName]}
case s := <-pm.hlsServerSet: case s := <-pm.hlsServerSet:
pm.hlsServer = s pm.hlsServer = s
@ -245,8 +245,8 @@ outer:
paths[name] = pa paths[name] = pa
} }
req.Res <- pathAPIPathsListRes{ req.res <- pathAPIPathsListRes{
Paths: paths, paths: paths,
} }
case <-pm.ctx.Done(): case <-pm.ctx.Done():
@ -332,52 +332,52 @@ func (pm *pathManager) onPathClose(pa *path) {
// onDescribe is called by a reader or publisher. // onDescribe is called by a reader or publisher.
func (pm *pathManager) onDescribe(req pathDescribeReq) pathDescribeRes { func (pm *pathManager) onDescribe(req pathDescribeReq) pathDescribeRes {
req.Res = make(chan pathDescribeRes) req.res = make(chan pathDescribeRes)
select { select {
case pm.describe <- req: case pm.describe <- req:
res := <-req.Res res := <-req.res
if res.Err != nil { if res.err != nil {
return res return res
} }
return res.Path.onDescribe(req) return res.path.onDescribe(req)
case <-pm.ctx.Done(): case <-pm.ctx.Done():
return pathDescribeRes{Err: fmt.Errorf("terminated")} return pathDescribeRes{err: fmt.Errorf("terminated")}
} }
} }
// onPublisherAnnounce is called by a publisher. // onPublisherAnnounce is called by a publisher.
func (pm *pathManager) onPublisherAnnounce(req pathPublisherAnnounceReq) pathPublisherAnnounceRes { func (pm *pathManager) onPublisherAnnounce(req pathPublisherAnnounceReq) pathPublisherAnnounceRes {
req.Res = make(chan pathPublisherAnnounceRes) req.res = make(chan pathPublisherAnnounceRes)
select { select {
case pm.publisherAnnounce <- req: case pm.publisherAnnounce <- req:
res := <-req.Res res := <-req.res
if res.Err != nil { if res.err != nil {
return res return res
} }
return res.Path.onPublisherAnnounce(req) return res.path.onPublisherAnnounce(req)
case <-pm.ctx.Done(): case <-pm.ctx.Done():
return pathPublisherAnnounceRes{Err: fmt.Errorf("terminated")} return pathPublisherAnnounceRes{err: fmt.Errorf("terminated")}
} }
} }
// onReaderSetupPlay is called by a reader. // onReaderSetupPlay is called by a reader.
func (pm *pathManager) onReaderSetupPlay(req pathReaderSetupPlayReq) pathReaderSetupPlayRes { func (pm *pathManager) onReaderSetupPlay(req pathReaderSetupPlayReq) pathReaderSetupPlayRes {
req.Res = make(chan pathReaderSetupPlayRes) req.res = make(chan pathReaderSetupPlayRes)
select { select {
case pm.readerSetupPlay <- req: case pm.readerSetupPlay <- req:
res := <-req.Res res := <-req.res
if res.Err != nil { if res.err != nil {
return res return res
} }
return res.Path.onReaderSetupPlay(req) return res.path.onReaderSetupPlay(req)
case <-pm.ctx.Done(): case <-pm.ctx.Done():
return pathReaderSetupPlayRes{Err: fmt.Errorf("terminated")} return pathReaderSetupPlayRes{err: fmt.Errorf("terminated")}
} }
} }
@ -391,22 +391,22 @@ func (pm *pathManager) onHLSServerSet(s pathManagerHLSServer) {
// onAPIPathsList is called by api. // onAPIPathsList is called by api.
func (pm *pathManager) onAPIPathsList(req pathAPIPathsListReq) pathAPIPathsListRes { func (pm *pathManager) onAPIPathsList(req pathAPIPathsListReq) pathAPIPathsListRes {
req.Res = make(chan pathAPIPathsListRes) req.res = make(chan pathAPIPathsListRes)
select { select {
case pm.apiPathsList <- req: case pm.apiPathsList <- req:
res := <-req.Res res := <-req.res
res.Data = &pathAPIPathsListData{ res.data = &pathAPIPathsListData{
Items: make(map[string]pathAPIPathsListItem), Items: make(map[string]pathAPIPathsListItem),
} }
for _, pa := range res.Paths { for _, pa := range res.paths {
pa.onAPIPathsList(pathAPIPathsListSubReq{Data: res.Data}) pa.onAPIPathsList(pathAPIPathsListSubReq{data: res.data})
} }
return res return res
case <-pm.ctx.Done(): case <-pm.ctx.Done():
return pathAPIPathsListRes{Err: fmt.Errorf("terminated")} return pathAPIPathsListRes{err: fmt.Errorf("terminated")}
} }
} }

58
internal/core/rtmp_conn.go

@ -220,9 +220,9 @@ func (c *rtmpConn) runRead(ctx context.Context) error {
pathName, query := pathNameAndQuery(c.conn.URL()) pathName, query := pathNameAndQuery(c.conn.URL())
res := c.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{ res := c.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{
Author: c, author: c,
PathName: pathName, pathName: pathName,
Authenticate: func( authenticate: func(
pathIPs []interface{}, pathIPs []interface{},
pathUser conf.Credential, pathUser conf.Credential,
pathPass conf.Credential) error { pathPass conf.Credential) error {
@ -230,19 +230,19 @@ func (c *rtmpConn) runRead(ctx context.Context) error {
}, },
}) })
if res.Err != nil { if res.err != nil {
if terr, ok := res.Err.(pathErrAuthCritical); ok { if terr, ok := res.err.(pathErrAuthCritical); ok {
// wait some seconds to stop brute force attacks // wait some seconds to stop brute force attacks
<-time.After(rtmpConnPauseAfterAuthError) <-time.After(rtmpConnPauseAfterAuthError)
return errors.New(terr.Message) return errors.New(terr.message)
} }
return res.Err return res.err
} }
c.path = res.Path c.path = res.path
defer func() { defer func() {
c.path.onReaderRemove(pathReaderRemoveReq{Author: c}) c.path.onReaderRemove(pathReaderRemoveReq{author: c})
}() }()
c.stateMutex.Lock() c.stateMutex.Lock()
@ -257,7 +257,7 @@ func (c *rtmpConn) runRead(ctx context.Context) error {
var audioClockRate int var audioClockRate int
var aacDecoder *rtpaac.Decoder var aacDecoder *rtpaac.Decoder
for i, t := range res.Stream.tracks() { for i, t := range res.stream.tracks() {
if t.IsH264() { if t.IsH264() {
if videoTrack != nil { if videoTrack != nil {
return fmt.Errorf("can't read track %d with RTMP: too many tracks", i+1) return fmt.Errorf("can't read track %d with RTMP: too many tracks", i+1)
@ -293,7 +293,7 @@ func (c *rtmpConn) runRead(ctx context.Context) error {
}() }()
c.path.onReaderPlay(pathReaderPlayReq{ c.path.onReaderPlay(pathReaderPlayReq{
Author: c, author: c,
}) })
if c.path.Conf().RunOnRead != "" { if c.path.Conf().RunOnRead != "" {
@ -465,9 +465,9 @@ func (c *rtmpConn) runPublish(ctx context.Context) error {
pathName, query := pathNameAndQuery(c.conn.URL()) pathName, query := pathNameAndQuery(c.conn.URL())
res := c.pathManager.onPublisherAnnounce(pathPublisherAnnounceReq{ res := c.pathManager.onPublisherAnnounce(pathPublisherAnnounceReq{
Author: c, author: c,
PathName: pathName, pathName: pathName,
Authenticate: func( authenticate: func(
pathIPs []interface{}, pathIPs []interface{},
pathUser conf.Credential, pathUser conf.Credential,
pathPass conf.Credential) error { pathPass conf.Credential) error {
@ -475,19 +475,19 @@ func (c *rtmpConn) runPublish(ctx context.Context) error {
}, },
}) })
if res.Err != nil { if res.err != nil {
if terr, ok := res.Err.(pathErrAuthCritical); ok { if terr, ok := res.err.(pathErrAuthCritical); ok {
// wait some seconds to stop brute force attacks // wait some seconds to stop brute force attacks
<-time.After(rtmpConnPauseAfterAuthError) <-time.After(rtmpConnPauseAfterAuthError)
return errors.New(terr.Message) return errors.New(terr.message)
} }
return res.Err return res.err
} }
c.path = res.Path c.path = res.path
defer func() { defer func() {
c.path.onPublisherRemove(pathPublisherRemoveReq{Author: c}) c.path.onPublisherRemove(pathPublisherRemoveReq{author: c})
}() }()
c.stateMutex.Lock() c.stateMutex.Lock()
@ -498,19 +498,19 @@ func (c *rtmpConn) runPublish(ctx context.Context) error {
c.conn.SetWriteDeadline(time.Time{}) c.conn.SetWriteDeadline(time.Time{})
rres := c.path.onPublisherRecord(pathPublisherRecordReq{ rres := c.path.onPublisherRecord(pathPublisherRecordReq{
Author: c, author: c,
Tracks: tracks, tracks: tracks,
}) })
if rres.Err != nil { if rres.err != nil {
return rres.Err return rres.err
} }
rtcpSenders := rtcpsenderset.New(tracks, rres.Stream.onPacketRTCP) rtcpSenders := rtcpsenderset.New(tracks, rres.stream.onPacketRTCP)
defer rtcpSenders.Close() defer rtcpSenders.Close()
onPacketRTP := func(trackID int, payload []byte) { onPacketRTP := func(trackID int, payload []byte) {
rtcpSenders.OnPacketRTP(trackID, payload) rtcpSenders.OnPacketRTP(trackID, payload)
rres.Stream.onPacketRTP(trackID, payload) rres.stream.onPacketRTP(trackID, payload)
} }
for { for {
@ -610,7 +610,7 @@ func (c *rtmpConn) authenticate(
action) action)
if err != nil { if err != nil {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: fmt.Sprintf("external authentication failed: %s", err), message: fmt.Sprintf("external authentication failed: %s", err),
} }
} }
} }
@ -619,7 +619,7 @@ func (c *rtmpConn) authenticate(
ip := c.ip() ip := c.ip()
if !ipEqualOrInRange(ip, pathIPs) { if !ipEqualOrInRange(ip, pathIPs) {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: fmt.Sprintf("IP '%s' not allowed", ip), message: fmt.Sprintf("IP '%s' not allowed", ip),
} }
} }
} }
@ -628,7 +628,7 @@ func (c *rtmpConn) authenticate(
if query.Get("user") != string(pathUser) || if query.Get("user") != string(pathUser) ||
query.Get("pass") != string(pathPass) { query.Get("pass") != string(pathPass) {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: "invalid credentials", message: "invalid credentials",
} }
} }
} }

32
internal/core/rtmp_server.go

@ -26,21 +26,21 @@ type rtmpServerAPIConnsListData struct {
} }
type rtmpServerAPIConnsListRes struct { type rtmpServerAPIConnsListRes struct {
Data *rtmpServerAPIConnsListData data *rtmpServerAPIConnsListData
Err error err error
} }
type rtmpServerAPIConnsListReq struct { type rtmpServerAPIConnsListReq struct {
Res chan rtmpServerAPIConnsListRes res chan rtmpServerAPIConnsListRes
} }
type rtmpServerAPIConnsKickRes struct { type rtmpServerAPIConnsKickRes struct {
Err error err error
} }
type rtmpServerAPIConnsKickReq struct { type rtmpServerAPIConnsKickReq struct {
ID string id string
Res chan rtmpServerAPIConnsKickRes res chan rtmpServerAPIConnsKickRes
} }
type rtmpServerParent interface { type rtmpServerParent interface {
@ -219,12 +219,12 @@ outer:
} }
} }
req.Res <- rtmpServerAPIConnsListRes{Data: data} req.res <- rtmpServerAPIConnsListRes{data: data}
case req := <-s.apiConnsKick: case req := <-s.apiConnsKick:
res := func() bool { res := func() bool {
for c := range s.conns { for c := range s.conns {
if c.ID() == req.ID { if c.ID() == req.id {
delete(s.conns, c) delete(s.conns, c)
c.close() c.close()
return true return true
@ -233,9 +233,9 @@ outer:
return false return false
}() }()
if res { if res {
req.Res <- rtmpServerAPIConnsKickRes{} req.res <- rtmpServerAPIConnsKickRes{}
} else { } else {
req.Res <- rtmpServerAPIConnsKickRes{fmt.Errorf("not found")} req.res <- rtmpServerAPIConnsKickRes{fmt.Errorf("not found")}
} }
case <-s.ctx.Done(): case <-s.ctx.Done():
@ -290,24 +290,24 @@ func (s *rtmpServer) onConnClose(c *rtmpConn) {
// onAPIConnsList is called by api. // onAPIConnsList is called by api.
func (s *rtmpServer) onAPIConnsList(req rtmpServerAPIConnsListReq) rtmpServerAPIConnsListRes { func (s *rtmpServer) onAPIConnsList(req rtmpServerAPIConnsListReq) rtmpServerAPIConnsListRes {
req.Res = make(chan rtmpServerAPIConnsListRes) req.res = make(chan rtmpServerAPIConnsListRes)
select { select {
case s.apiConnsList <- req: case s.apiConnsList <- req:
return <-req.Res return <-req.res
case <-s.ctx.Done(): case <-s.ctx.Done():
return rtmpServerAPIConnsListRes{Err: fmt.Errorf("terminated")} return rtmpServerAPIConnsListRes{err: fmt.Errorf("terminated")}
} }
} }
// onAPIConnsKick is called by api. // onAPIConnsKick is called by api.
func (s *rtmpServer) onAPIConnsKick(req rtmpServerAPIConnsKickReq) rtmpServerAPIConnsKickRes { func (s *rtmpServer) onAPIConnsKick(req rtmpServerAPIConnsKickReq) rtmpServerAPIConnsKickRes {
req.Res = make(chan rtmpServerAPIConnsKickRes) req.res = make(chan rtmpServerAPIConnsKickRes)
select { select {
case s.apiConnsKick <- req: case s.apiConnsKick <- req:
return <-req.Res return <-req.res
case <-s.ctx.Done(): case <-s.ctx.Done():
return rtmpServerAPIConnsKickRes{Err: fmt.Errorf("terminated")} return rtmpServerAPIConnsKickRes{err: fmt.Errorf("terminated")}
} }
} }

16
internal/core/rtmp_source.go

@ -25,7 +25,7 @@ const (
type rtmpSourceParent interface { type rtmpSourceParent interface {
log(logger.Level, string, ...interface{}) log(logger.Level, string, ...interface{})
onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes
OnSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) onSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq)
} }
type rtmpSource struct { type rtmpSource struct {
@ -149,25 +149,25 @@ func (s *rtmpSource) runInner() bool {
} }
res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{ res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{
Source: s, source: s,
Tracks: tracks, tracks: tracks,
}) })
if res.Err != nil { if res.err != nil {
return res.Err return res.err
} }
s.log(logger.Info, "ready") s.log(logger.Info, "ready")
defer func() { defer func() {
s.parent.OnSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{Source: s}) s.parent.onSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{source: s})
}() }()
rtcpSenders := rtcpsenderset.New(tracks, res.Stream.onPacketRTCP) rtcpSenders := rtcpsenderset.New(tracks, res.stream.onPacketRTCP)
defer rtcpSenders.Close() defer rtcpSenders.Close()
onPacketRTP := func(trackID int, payload []byte) { onPacketRTP := func(trackID int, payload []byte) {
rtcpSenders.OnPacketRTP(trackID, payload) rtcpSenders.OnPacketRTP(trackID, payload)
res.Stream.onPacketRTP(trackID, payload) res.stream.onPacketRTP(trackID, payload)
} }
for { for {

44
internal/core/rtsp_conn.go

@ -138,8 +138,8 @@ func (c *rtspConn) authenticate(
// therefore we must allow up to 3 failures // therefore we must allow up to 3 failures
if c.authFailures > 3 { if c.authFailures > 3 {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: "unauthorized: " + err.Error(), message: "unauthorized: " + err.Error(),
Response: &base.Response{ response: &base.Response{
StatusCode: base.StatusUnauthorized, StatusCode: base.StatusUnauthorized,
}, },
} }
@ -147,8 +147,8 @@ func (c *rtspConn) authenticate(
v := "IPCAM" v := "IPCAM"
return pathErrAuthNotCritical{ return pathErrAuthNotCritical{
Message: "unauthorized: " + err.Error(), message: "unauthorized: " + err.Error(),
Response: &base.Response{ response: &base.Response{
StatusCode: base.StatusUnauthorized, StatusCode: base.StatusUnauthorized,
Header: base.Header{ Header: base.Header{
"WWW-Authenticate": headers.Authenticate{ "WWW-Authenticate": headers.Authenticate{
@ -165,8 +165,8 @@ func (c *rtspConn) authenticate(
ip := c.ip() ip := c.ip()
if !ipEqualOrInRange(ip, pathIPs) { if !ipEqualOrInRange(ip, pathIPs) {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: fmt.Sprintf("IP '%s' not allowed", ip), message: fmt.Sprintf("IP '%s' not allowed", ip),
Response: &base.Response{ response: &base.Response{
StatusCode: base.StatusUnauthorized, StatusCode: base.StatusUnauthorized,
}, },
} }
@ -193,15 +193,15 @@ func (c *rtspConn) authenticate(
// therefore we must allow up to 3 failures // therefore we must allow up to 3 failures
if c.authFailures > 3 { if c.authFailures > 3 {
return pathErrAuthCritical{ return pathErrAuthCritical{
Message: "unauthorized: " + err.Error(), message: "unauthorized: " + err.Error(),
Response: &base.Response{ response: &base.Response{
StatusCode: base.StatusUnauthorized, StatusCode: base.StatusUnauthorized,
}, },
} }
} }
return pathErrAuthNotCritical{ return pathErrAuthNotCritical{
Response: &base.Response{ response: &base.Response{
StatusCode: base.StatusUnauthorized, StatusCode: base.StatusUnauthorized,
Header: base.Header{ Header: base.Header{
"WWW-Authenticate": c.authValidator.Header(), "WWW-Authenticate": c.authValidator.Header(),
@ -241,9 +241,9 @@ func (c *rtspConn) OnResponse(res *base.Response) {
func (c *rtspConn) onDescribe(ctx *gortsplib.ServerHandlerOnDescribeCtx, func (c *rtspConn) onDescribe(ctx *gortsplib.ServerHandlerOnDescribeCtx,
) (*base.Response, *gortsplib.ServerStream, error) { ) (*base.Response, *gortsplib.ServerStream, error) {
res := c.pathManager.onDescribe(pathDescribeReq{ res := c.pathManager.onDescribe(pathDescribeReq{
PathName: ctx.Path, pathName: ctx.Path,
URL: ctx.Req.URL, url: ctx.Req.URL,
Authenticate: func( authenticate: func(
pathIPs []interface{}, pathIPs []interface{},
pathUser conf.Credential, pathUser conf.Credential,
pathPass conf.Credential) error { pathPass conf.Credential) error {
@ -251,40 +251,40 @@ func (c *rtspConn) onDescribe(ctx *gortsplib.ServerHandlerOnDescribeCtx,
}, },
}) })
if res.Err != nil { if res.err != nil {
switch terr := res.Err.(type) { switch terr := res.err.(type) {
case pathErrAuthNotCritical: case pathErrAuthNotCritical:
c.log(logger.Debug, "non-critical authentication error: %s", terr.Message) c.log(logger.Debug, "non-critical authentication error: %s", terr.message)
return terr.Response, nil, nil return terr.response, nil, nil
case pathErrAuthCritical: case pathErrAuthCritical:
// wait some seconds to stop brute force attacks // wait some seconds to stop brute force attacks
<-time.After(rtspConnPauseAfterAuthError) <-time.After(rtspConnPauseAfterAuthError)
return terr.Response, nil, errors.New(terr.Message) return terr.response, nil, errors.New(terr.message)
case pathErrNoOnePublishing: case pathErrNoOnePublishing:
return &base.Response{ return &base.Response{
StatusCode: base.StatusNotFound, StatusCode: base.StatusNotFound,
}, nil, res.Err }, nil, res.err
default: default:
return &base.Response{ return &base.Response{
StatusCode: base.StatusBadRequest, StatusCode: base.StatusBadRequest,
}, nil, res.Err }, nil, res.err
} }
} }
if res.Redirect != "" { if res.redirect != "" {
return &base.Response{ return &base.Response{
StatusCode: base.StatusMovedPermanently, StatusCode: base.StatusMovedPermanently,
Header: base.Header{ Header: base.Header{
"Location": base.HeaderValue{res.Redirect}, "Location": base.HeaderValue{res.redirect},
}, },
}, nil, nil }, nil, nil
} }
return &base.Response{ return &base.Response{
StatusCode: base.StatusOK, StatusCode: base.StatusOK,
}, res.Stream.rtspStream, nil }, res.stream.rtspStream, nil
} }

18
internal/core/rtsp_server.go

@ -31,18 +31,18 @@ type rtspServerAPISessionsListData struct {
} }
type rtspServerAPISessionsListRes struct { type rtspServerAPISessionsListRes struct {
Data *rtspServerAPISessionsListData data *rtspServerAPISessionsListData
Err error err error
} }
type rtspServerAPISessionsListReq struct{} type rtspServerAPISessionsListReq struct{}
type rtspServerAPISessionsKickRes struct { type rtspServerAPISessionsKickRes struct {
Err error err error
} }
type rtspServerAPISessionsKickReq struct { type rtspServerAPISessionsKickReq struct {
ID string id string
} }
type rtspServerParent interface { type rtspServerParent interface {
@ -405,7 +405,7 @@ func (s *rtspServer) OnPacketRTCP(ctx *gortsplib.ServerHandlerOnPacketRTCPCtx) {
func (s *rtspServer) onAPISessionsList(req rtspServerAPISessionsListReq) rtspServerAPISessionsListRes { func (s *rtspServer) onAPISessionsList(req rtspServerAPISessionsListReq) rtspServerAPISessionsListRes {
select { select {
case <-s.ctx.Done(): case <-s.ctx.Done():
return rtspServerAPISessionsListRes{Err: fmt.Errorf("terminated")} return rtspServerAPISessionsListRes{err: fmt.Errorf("terminated")}
default: default:
} }
@ -434,14 +434,14 @@ func (s *rtspServer) onAPISessionsList(req rtspServerAPISessionsListReq) rtspSer
} }
} }
return rtspServerAPISessionsListRes{Data: data} return rtspServerAPISessionsListRes{data: data}
} }
// onAPISessionsKick is called by api. // onAPISessionsKick is called by api.
func (s *rtspServer) onAPISessionsKick(req rtspServerAPISessionsKickReq) rtspServerAPISessionsKickRes { func (s *rtspServer) onAPISessionsKick(req rtspServerAPISessionsKickReq) rtspServerAPISessionsKickRes {
select { select {
case <-s.ctx.Done(): case <-s.ctx.Done():
return rtspServerAPISessionsKickRes{Err: fmt.Errorf("terminated")} return rtspServerAPISessionsKickRes{err: fmt.Errorf("terminated")}
default: default:
} }
@ -449,7 +449,7 @@ func (s *rtspServer) onAPISessionsKick(req rtspServerAPISessionsKickReq) rtspSer
defer s.mutex.RUnlock() defer s.mutex.RUnlock()
for key, se := range s.sessions { for key, se := range s.sessions {
if se.ID() == req.ID { if se.ID() == req.id {
se.close() se.close()
delete(s.sessions, key) delete(s.sessions, key)
se.onClose(liberrors.ErrServerTerminated{}) se.onClose(liberrors.ErrServerTerminated{})
@ -457,5 +457,5 @@ func (s *rtspServer) onAPISessionsKick(req rtspServerAPISessionsKickReq) rtspSer
} }
} }
return rtspServerAPISessionsKickRes{Err: fmt.Errorf("not found")} return rtspServerAPISessionsKickRes{err: fmt.Errorf("not found")}
} }

68
internal/core/rtsp_session.go

@ -112,11 +112,11 @@ func (s *rtspSession) onClose(err error) {
switch s.ss.State() { switch s.ss.State() {
case gortsplib.ServerSessionStatePreRead, gortsplib.ServerSessionStateRead: case gortsplib.ServerSessionStatePreRead, gortsplib.ServerSessionStateRead:
s.path.onReaderRemove(pathReaderRemoveReq{Author: s}) s.path.onReaderRemove(pathReaderRemoveReq{author: s})
s.path = nil s.path = nil
case gortsplib.ServerSessionStatePrePublish, gortsplib.ServerSessionStatePublish: case gortsplib.ServerSessionStatePrePublish, gortsplib.ServerSessionStatePublish:
s.path.onPublisherRemove(pathPublisherRemoveReq{Author: s}) s.path.onPublisherRemove(pathPublisherRemoveReq{author: s})
s.path = nil s.path = nil
} }
@ -155,9 +155,9 @@ func (s *rtspSession) onAnnounce(c *rtspConn, ctx *gortsplib.ServerHandlerOnAnno
} }
res := s.pathManager.onPublisherAnnounce(pathPublisherAnnounceReq{ res := s.pathManager.onPublisherAnnounce(pathPublisherAnnounceReq{
Author: s, author: s,
PathName: ctx.Path, pathName: ctx.Path,
Authenticate: func( authenticate: func(
pathIPs []interface{}, pathIPs []interface{},
pathUser conf.Credential, pathUser conf.Credential,
pathPass conf.Credential) error { pathPass conf.Credential) error {
@ -165,26 +165,26 @@ func (s *rtspSession) onAnnounce(c *rtspConn, ctx *gortsplib.ServerHandlerOnAnno
}, },
}) })
if res.Err != nil { if res.err != nil {
switch terr := res.Err.(type) { switch terr := res.err.(type) {
case pathErrAuthNotCritical: case pathErrAuthNotCritical:
s.log(logger.Debug, "non-critical authentication error: %s", terr.Message) s.log(logger.Debug, "non-critical authentication error: %s", terr.message)
return terr.Response, nil return terr.response, nil
case pathErrAuthCritical: case pathErrAuthCritical:
// wait some seconds to stop brute force attacks // wait some seconds to stop brute force attacks
<-time.After(pauseAfterAuthError) <-time.After(pauseAfterAuthError)
return terr.Response, errors.New(terr.Message) return terr.response, errors.New(terr.message)
default: default:
return &base.Response{ return &base.Response{
StatusCode: base.StatusBadRequest, StatusCode: base.StatusBadRequest,
}, res.Err }, res.err
} }
} }
s.path = res.Path s.path = res.path
s.announcedTracks = ctx.Tracks s.announcedTracks = ctx.Tracks
s.stateMutex.Lock() s.stateMutex.Lock()
@ -214,9 +214,9 @@ func (s *rtspSession) onSetup(c *rtspConn, ctx *gortsplib.ServerHandlerOnSetupCt
switch s.ss.State() { switch s.ss.State() {
case gortsplib.ServerSessionStateInitial, gortsplib.ServerSessionStatePreRead: // play case gortsplib.ServerSessionStateInitial, gortsplib.ServerSessionStatePreRead: // play
res := s.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{ res := s.pathManager.onReaderSetupPlay(pathReaderSetupPlayReq{
Author: s, author: s,
PathName: ctx.Path, pathName: ctx.Path,
Authenticate: func( authenticate: func(
pathIPs []interface{}, pathIPs []interface{},
pathUser conf.Credential, pathUser conf.Credential,
pathPass conf.Credential) error { pathPass conf.Credential) error {
@ -224,33 +224,33 @@ func (s *rtspSession) onSetup(c *rtspConn, ctx *gortsplib.ServerHandlerOnSetupCt
}, },
}) })
if res.Err != nil { if res.err != nil {
switch terr := res.Err.(type) { switch terr := res.err.(type) {
case pathErrAuthNotCritical: case pathErrAuthNotCritical:
s.log(logger.Debug, "non-critical authentication error: %s", terr.Message) s.log(logger.Debug, "non-critical authentication error: %s", terr.message)
return terr.Response, nil, nil return terr.response, nil, nil
case pathErrAuthCritical: case pathErrAuthCritical:
// wait some seconds to stop brute force attacks // wait some seconds to stop brute force attacks
<-time.After(pauseAfterAuthError) <-time.After(pauseAfterAuthError)
return terr.Response, nil, errors.New(terr.Message) return terr.response, nil, errors.New(terr.message)
case pathErrNoOnePublishing: case pathErrNoOnePublishing:
return &base.Response{ return &base.Response{
StatusCode: base.StatusNotFound, StatusCode: base.StatusNotFound,
}, nil, res.Err }, nil, res.err
default: default:
return &base.Response{ return &base.Response{
StatusCode: base.StatusBadRequest, StatusCode: base.StatusBadRequest,
}, nil, res.Err }, nil, res.err
} }
} }
s.path = res.Path s.path = res.path
if ctx.TrackID >= len(res.Stream.tracks()) { if ctx.TrackID >= len(res.stream.tracks()) {
return &base.Response{ return &base.Response{
StatusCode: base.StatusBadRequest, StatusCode: base.StatusBadRequest,
}, nil, fmt.Errorf("track %d does not exist", ctx.TrackID) }, nil, fmt.Errorf("track %d does not exist", ctx.TrackID)
@ -259,7 +259,7 @@ func (s *rtspSession) onSetup(c *rtspConn, ctx *gortsplib.ServerHandlerOnSetupCt
if s.setuppedTracks == nil { if s.setuppedTracks == nil {
s.setuppedTracks = make(map[int]*gortsplib.Track) s.setuppedTracks = make(map[int]*gortsplib.Track)
} }
s.setuppedTracks[ctx.TrackID] = res.Stream.tracks()[ctx.TrackID] s.setuppedTracks[ctx.TrackID] = res.stream.tracks()[ctx.TrackID]
s.stateMutex.Lock() s.stateMutex.Lock()
s.state = gortsplib.ServerSessionStatePreRead s.state = gortsplib.ServerSessionStatePreRead
@ -267,7 +267,7 @@ func (s *rtspSession) onSetup(c *rtspConn, ctx *gortsplib.ServerHandlerOnSetupCt
return &base.Response{ return &base.Response{
StatusCode: base.StatusOK, StatusCode: base.StatusOK,
}, res.Stream.rtspStream, nil }, res.stream.rtspStream, nil
default: // record default: // record
return &base.Response{ return &base.Response{
@ -281,7 +281,7 @@ func (s *rtspSession) onPlay(ctx *gortsplib.ServerHandlerOnPlayCtx) (*base.Respo
h := make(base.Header) h := make(base.Header)
if s.ss.State() == gortsplib.ServerSessionStatePreRead { if s.ss.State() == gortsplib.ServerSessionStatePreRead {
s.path.onReaderPlay(pathReaderPlayReq{Author: s}) s.path.onReaderPlay(pathReaderPlayReq{author: s})
if s.path.Conf().RunOnRead != "" { if s.path.Conf().RunOnRead != "" {
s.log(logger.Info, "runOnRead command started") s.log(logger.Info, "runOnRead command started")
@ -309,16 +309,16 @@ func (s *rtspSession) onPlay(ctx *gortsplib.ServerHandlerOnPlayCtx) (*base.Respo
// onRecord is called by rtspServer. // onRecord is called by rtspServer.
func (s *rtspSession) onRecord(ctx *gortsplib.ServerHandlerOnRecordCtx) (*base.Response, error) { func (s *rtspSession) onRecord(ctx *gortsplib.ServerHandlerOnRecordCtx) (*base.Response, error) {
res := s.path.onPublisherRecord(pathPublisherRecordReq{ res := s.path.onPublisherRecord(pathPublisherRecordReq{
Author: s, author: s,
Tracks: s.announcedTracks, tracks: s.announcedTracks,
}) })
if res.Err != nil { if res.err != nil {
return &base.Response{ return &base.Response{
StatusCode: base.StatusBadRequest, StatusCode: base.StatusBadRequest,
}, res.Err }, res.err
} }
s.stream = res.Stream s.stream = res.stream
s.stateMutex.Lock() s.stateMutex.Lock()
s.state = gortsplib.ServerSessionStatePublish s.state = gortsplib.ServerSessionStatePublish
@ -338,14 +338,14 @@ func (s *rtspSession) onPause(ctx *gortsplib.ServerHandlerOnPauseCtx) (*base.Res
s.onReadCmd.Close() s.onReadCmd.Close()
} }
s.path.onReaderPause(pathReaderPauseReq{Author: s}) s.path.onReaderPause(pathReaderPauseReq{author: s})
s.stateMutex.Lock() s.stateMutex.Lock()
s.state = gortsplib.ServerSessionStatePreRead s.state = gortsplib.ServerSessionStatePreRead
s.stateMutex.Unlock() s.stateMutex.Unlock()
case gortsplib.ServerSessionStatePublish: case gortsplib.ServerSessionStatePublish:
s.path.onPublisherPause(pathPublisherPauseReq{Author: s}) s.path.onPublisherPause(pathPublisherPauseReq{author: s})
s.stateMutex.Lock() s.stateMutex.Lock()
s.state = gortsplib.ServerSessionStatePrePublish s.state = gortsplib.ServerSessionStatePrePublish

28
internal/core/rtsp_source.go

@ -27,7 +27,7 @@ const (
type rtspSourceParent interface { type rtspSourceParent interface {
log(logger.Level, string, ...interface{}) log(logger.Level, string, ...interface{})
onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes onSourceStaticSetReady(req pathSourceStaticSetReadyReq) pathSourceStaticSetReadyRes
OnSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) onSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq)
} }
type rtspSource struct { type rtspSource struct {
@ -195,25 +195,25 @@ func (s *rtspSource) runInner() bool {
} }
res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{ res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{
Source: s, source: s,
Tracks: c.Tracks(), tracks: c.Tracks(),
}) })
if res.Err != nil { if res.err != nil {
return res.Err return res.err
} }
s.log(logger.Info, "ready") s.log(logger.Info, "ready")
defer func() { defer func() {
s.parent.OnSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{Source: s}) s.parent.onSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{source: s})
}() }()
c.OnPacketRTP = func(trackID int, payload []byte) { c.OnPacketRTP = func(trackID int, payload []byte) {
res.Stream.onPacketRTP(trackID, payload) res.stream.onPacketRTP(trackID, payload)
} }
c.OnPacketRTCP = func(trackID int, payload []byte) { c.OnPacketRTCP = func(trackID int, payload []byte) {
res.Stream.onPacketRTCP(trackID, payload) res.stream.onPacketRTCP(trackID, payload)
} }
_, err = c.Play(nil) _, err = c.Play(nil)
@ -358,23 +358,23 @@ func (s *rtspSource) handleMissingH264Params(c *gortsplib.Client, tracks gortspl
tracks[h264TrackID] = track tracks[h264TrackID] = track
res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{ res := s.parent.onSourceStaticSetReady(pathSourceStaticSetReadyReq{
Source: s, source: s,
Tracks: tracks, tracks: tracks,
}) })
if res.Err != nil { if res.err != nil {
return res.Err return res.err
} }
func() { func() {
streamMutex.Lock() streamMutex.Lock()
defer streamMutex.Unlock() defer streamMutex.Unlock()
stream = res.Stream stream = res.stream
}() }()
s.log(logger.Info, "ready") s.log(logger.Info, "ready")
defer func() { defer func() {
s.parent.OnSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{Source: s}) s.parent.onSourceStaticSetNotReady(pathSourceStaticSetNotReadyReq{source: s})
}() }()
} }

Loading…
Cancel
Save