From f07886db5ff076d5670d4e79d6380987faa3d271 Mon Sep 17 00:00:00 2001 From: Alessandro Ros Date: Sat, 9 Sep 2023 23:37:56 +0200 Subject: [PATCH] print the reason why a source is started or stopped (#2322) --- internal/core/hls_muxer.go | 2 +- internal/core/path.go | 34 +++++++++++++++++---------------- internal/core/rtmp_conn.go | 2 +- internal/core/rtsp_conn.go | 2 +- internal/core/rtsp_session.go | 2 +- internal/core/source_static.go | 20 +++++++++++-------- internal/core/srt_conn.go | 2 +- internal/core/webrtc_session.go | 2 +- 8 files changed, 36 insertions(+), 30 deletions(-) diff --git a/internal/core/hls_muxer.go b/internal/core/hls_muxer.go index 58974524..22443a8a 100644 --- a/internal/core/hls_muxer.go +++ b/internal/core/hls_muxer.go @@ -230,7 +230,7 @@ func (m *hlsMuxer) run() { m.parent.closeMuxer(m) - m.Log(logger.Info, "destroyed (%v)", err) + m.Log(logger.Info, "destroyed: %v", err) } func (m *hlsMuxer) clearQueuedRequests() { diff --git a/internal/core/path.go b/internal/core/path.go index 226c911f..77e85c1f 100644 --- a/internal/core/path.go +++ b/internal/core/path.go @@ -314,7 +314,7 @@ func (pa *path) run() { pa) if !pa.conf.SourceOnDemand { - pa.source.(*sourceStatic).start() + pa.source.(*sourceStatic).start(false) } } @@ -362,7 +362,9 @@ func (pa *path) run() { if pa.source != nil { if source, ok := pa.source.(*sourceStatic); ok { - source.close() + if !pa.conf.SourceOnDemand || pa.onDemandStaticSourceState != pathOnDemandStateInitial { + source.close("path is closing") + } } else if source, ok := pa.source.(publisher); ok { source.close() } @@ -373,7 +375,7 @@ func (pa *path) run() { pa.Log(logger.Info, "runOnDemand command stopped") } - pa.Log(logger.Debug, "destroyed (%v)", err) + pa.Log(logger.Debug, "destroyed: %v", err) } func (pa *path) runInner() error { @@ -477,12 +479,12 @@ func (pa *path) doOnDemandStaticSourceReadyTimer() { } pa.readerAddRequestsOnHold = nil - pa.onDemandStaticSourceStop() + pa.onDemandStaticSourceStop("timed out") } func (pa *path) doOnDemandStaticSourceCloseTimer() { pa.setNotReady() - pa.onDemandStaticSourceStop() + pa.onDemandStaticSourceStop("not needed by anyone") } func (pa *path) doOnDemandPublisherReadyTimer() { @@ -496,11 +498,11 @@ func (pa *path) doOnDemandPublisherReadyTimer() { } pa.readerAddRequestsOnHold = nil - pa.onDemandStopPublisher() + pa.onDemandPublisherStop("timed out") } func (pa *path) doOnDemandPublisherCloseTimer() { - pa.onDemandStopPublisher() + pa.onDemandPublisherStop("not needed by anyone") } func (pa *path) doReloadConf(newConf *conf.PathConf) { @@ -550,7 +552,7 @@ func (pa *path) doSourceStaticSetNotReady(req pathSourceStaticSetNotReadyReq) { close(req.res) if pa.conf.HasOnDemandStaticSource() && pa.onDemandStaticSourceState != pathOnDemandStateInitial { - pa.onDemandStaticSourceStop() + pa.onDemandStaticSourceStop("an error occurred") } } @@ -579,7 +581,7 @@ func (pa *path) doDescribe(req pathDescribeReq) { if pa.conf.HasOnDemandPublisher() { if pa.onDemandPublisherState == pathOnDemandStateInitial { - pa.onDemandStartPublisher() + pa.onDemandPublisherStart() } pa.describeRequestsOnHold = append(pa.describeRequestsOnHold, req) return @@ -697,7 +699,7 @@ func (pa *path) doAddReader(req pathAddReaderReq) { if pa.conf.HasOnDemandPublisher() { if pa.onDemandPublisherState == pathOnDemandStateInitial { - pa.onDemandStartPublisher() + pa.onDemandPublisherStart() } pa.readerAddRequestsOnHold = append(pa.readerAddRequestsOnHold, req) return @@ -790,7 +792,7 @@ func (pa *path) externalCmdEnv() externalcmd.Environment { } func (pa *path) onDemandStaticSourceStart() { - pa.source.(*sourceStatic).start() + pa.source.(*sourceStatic).start(true) pa.onDemandStaticSourceReadyTimer.Stop() pa.onDemandStaticSourceReadyTimer = time.NewTimer(time.Duration(pa.conf.SourceOnDemandStartTimeout)) @@ -805,7 +807,7 @@ func (pa *path) onDemandStaticSourceScheduleClose() { pa.onDemandStaticSourceState = pathOnDemandStateClosing } -func (pa *path) onDemandStaticSourceStop() { +func (pa *path) onDemandStaticSourceStop(reason string) { if pa.onDemandStaticSourceState == pathOnDemandStateClosing { pa.onDemandStaticSourceCloseTimer.Stop() pa.onDemandStaticSourceCloseTimer = newEmptyTimer() @@ -813,10 +815,10 @@ func (pa *path) onDemandStaticSourceStop() { pa.onDemandStaticSourceState = pathOnDemandStateInitial - pa.source.(*sourceStatic).stop() + pa.source.(*sourceStatic).stop(reason) } -func (pa *path) onDemandStartPublisher() { +func (pa *path) onDemandPublisherStart() { pa.Log(logger.Info, "runOnDemand command started") pa.onDemandCmd = externalcmd.NewCmd( pa.externalCmdPool, @@ -840,7 +842,7 @@ func (pa *path) onDemandPublisherScheduleClose() { pa.onDemandPublisherState = pathOnDemandStateClosing } -func (pa *path) onDemandStopPublisher() { +func (pa *path) onDemandPublisherStop(reason string) { if pa.source != nil { pa.source.(publisher).close() pa.executeRemovePublisher() @@ -856,7 +858,7 @@ func (pa *path) onDemandStopPublisher() { if pa.onDemandCmd != nil { pa.onDemandCmd.Close() pa.onDemandCmd = nil - pa.Log(logger.Info, "runOnDemand command stopped") + pa.Log(logger.Info, "runOnDemand command stopped: %s", reason) } } diff --git a/internal/core/rtmp_conn.go b/internal/core/rtmp_conn.go index 968fed39..78f81c3b 100644 --- a/internal/core/rtmp_conn.go +++ b/internal/core/rtmp_conn.go @@ -170,7 +170,7 @@ func (c *rtmpConn) run() { c.parent.closeConn(c) - c.Log(logger.Info, "closed (%v)", err) + c.Log(logger.Info, "closed: %v", err) } func (c *rtmpConn) runInner() error { diff --git a/internal/core/rtsp_conn.go b/internal/core/rtsp_conn.go index f8f96e64..44d361ca 100644 --- a/internal/core/rtsp_conn.go +++ b/internal/core/rtsp_conn.go @@ -110,7 +110,7 @@ func (c *rtspConn) ip() net.IP { // onClose is called by rtspServer. func (c *rtspConn) onClose(err error) { - c.Log(logger.Info, "closed (%v)", err) + c.Log(logger.Info, "closed: %v", err) if c.onConnectCmd != nil { c.onConnectCmd.Close() diff --git a/internal/core/rtsp_session.go b/internal/core/rtsp_session.go index 744a7a69..7e94bd05 100644 --- a/internal/core/rtsp_session.go +++ b/internal/core/rtsp_session.go @@ -117,7 +117,7 @@ func (s *rtspSession) onClose(err error) { s.path = nil s.stream = nil - s.Log(logger.Info, "destroyed (%v)", err) + s.Log(logger.Info, "destroyed: %v", err) } // onAnnounce is called by rtspServer. diff --git a/internal/core/source_static.go b/internal/core/source_static.go index 8de1da01..540df6af 100644 --- a/internal/core/source_static.go +++ b/internal/core/source_static.go @@ -105,19 +105,23 @@ func newSourceStatic( return s } -func (s *sourceStatic) close() { - if s.running { - s.stop() - } +func (s *sourceStatic) close(reason string) { + s.stop(reason) } -func (s *sourceStatic) start() { +func (s *sourceStatic) start(onDemand bool) { if s.running { panic("should not happen") } s.running = true - s.impl.Log(logger.Info, "started") + s.impl.Log(logger.Info, "started%s", + func() string { + if onDemand { + return " on demand" + } + return "" + }()) s.ctx, s.ctxCancel = context.WithCancel(context.Background()) s.done = make(chan struct{}) @@ -125,13 +129,13 @@ func (s *sourceStatic) start() { go s.run() } -func (s *sourceStatic) stop() { +func (s *sourceStatic) stop(reason string) { if !s.running { panic("should not happen") } s.running = false - s.impl.Log(logger.Info, "stopped") + s.impl.Log(logger.Info, "stopped: %s", reason) s.ctxCancel() diff --git a/internal/core/srt_conn.go b/internal/core/srt_conn.go index 257c9db3..d4dc3e7d 100644 --- a/internal/core/srt_conn.go +++ b/internal/core/srt_conn.go @@ -132,7 +132,7 @@ func (c *srtConn) run() { c.parent.closeConn(c) - c.Log(logger.Info, "closed (%v)", err) + c.Log(logger.Info, "closed: %v", err) } func (c *srtConn) runInner() error { diff --git a/internal/core/webrtc_session.go b/internal/core/webrtc_session.go index 3fd1b951..0cb7842b 100644 --- a/internal/core/webrtc_session.go +++ b/internal/core/webrtc_session.go @@ -248,7 +248,7 @@ func (s *webRTCSession) run() { s.parent.closeSession(s) - s.Log(logger.Info, "closed (%v)", err) + s.Log(logger.Info, "closed: %v", err) } func (s *webRTCSession) runInner() error {