From 8909e35a1706673aee000114f097f1819e559295 Mon Sep 17 00:00:00 2001 From: Alessandro Ros Date: Sat, 18 Jul 2026 17:41:38 +0200 Subject: [PATCH] add runOnOnline / runOnOffline hooks (#5399) (#5956) These are triggered and a stream is online (i.e. not just provided by an offline segment). --- api/openapi.yaml | 6 +++ docs/2-features/21-hooks.md | 40 +++++++++++++++++-- internal/conf/path.go | 3 ++ internal/core/path.go | 44 +++++++++++++++++---- internal/core/path_test.go | 76 +++++++++++++++++++++++++++++++++++++ internal/hooks/on_online.go | 69 +++++++++++++++++++++++++++++++++ mediamtx.yml | 21 +++++++++- 7 files changed, 247 insertions(+), 12 deletions(-) create mode 100644 internal/hooks/on_online.go diff --git a/api/openapi.yaml b/api/openapi.yaml index 47681ba6..efecc6e4 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -1139,6 +1139,12 @@ components: type: boolean runOnNotReady: type: string + runOnOffline: + type: string + runOnOnline: + type: string + runOnOnlineRestart: + type: boolean runOnRead: type: string runOnReadRestart: diff --git a/docs/2-features/21-hooks.md b/docs/2-features/21-hooks.md index adb21c25..581e782b 100644 --- a/docs/2-features/21-hooks.md +++ b/docs/2-features/21-hooks.md @@ -37,6 +37,7 @@ paths: mypath: # Command to run when this path is initialized. # This can be used to publish a stream when the server is launched. + # This is terminated with SIGINT when the program closes. # The following environment variables are available: # * MTX_PATH: path name # * RTSP_PORT: RTSP server port @@ -55,6 +56,7 @@ paths: pathDefaults: # Command to run when this path is requested by a reader # and no one is publishing to this path yet. + # This can be used to publish a stream on demand. # This is terminated with SIGINT when there are no readers anymore. # The following environment variables are available: # * MTX_PATH: path name @@ -80,12 +82,11 @@ pathDefaults: ## runOnReady -`runOnReady` allows to run a command when a stream is ready to be read: +`runOnReady` allows to run a command when a stream is available to be read: ```yml pathDefaults: - # Command to run when the stream is ready to be read, whenever it is - # published by a client or pulled from a server / camera. + # Command to run when the stream is available to be read. # This is terminated with SIGINT when the stream is not ready anymore. # The following environment variables are available: # * MTX_PATH: path name @@ -111,6 +112,39 @@ pathDefaults: runOnNotReady: curl http://my-custom-server/webhook?path=$MTX_PATH&source_type=$MTX_SOURCE_TYPE&source_id=$MTX_SOURCE_ID ``` +## runOnOnline + +`runOnOnline` allows to run a command when a stream is online, which means that the stream is available and provided by a online source (not an offline segment): + +```yml +pathDefaults: + # Command to run when the stream is online, which means + # that the stream is available and provided by a online source (not an offline segment). + # This is terminated with SIGINT when the stream is not online anymore. + # The following environment variables are available: + # * MTX_PATH: path name + # * MTX_QUERY: query parameters (passed by publisher) (url-encoded) + # * MTX_SOURCE_TYPE: source type + # * MTX_SOURCE_ID: source ID + # * RTSP_PORT: RTSP server port + # * G1, G2, ...: regular expression groups, if path name is + # a regular expression. + runOnOnline: curl http://my-custom-server/webhook?path=$MTX_PATH&source_type=$MTX_SOURCE_TYPE&source_id=$MTX_SOURCE_ID + # Restart the command if it exits. + runOnOnlineRestart: no +``` + +## runOnOffline + +`runOnOffline` allows to run a command when a stream is not online anymore: + +```yml +pathDefaults: + # Command to run when the stream is not online anymore. + # Environment variables are the same as runOnOnline. + runOnOffline: curl http://my-custom-server/webhook?path=$MTX_PATH&source_type=$MTX_SOURCE_TYPE&source_id=$MTX_SOURCE_ID +``` + ## runOnRead `runOnRead` allows to run a command when a client starts reading: diff --git a/internal/conf/path.go b/internal/conf/path.go index b53510c9..538329ec 100644 --- a/internal/conf/path.go +++ b/internal/conf/path.go @@ -344,6 +344,9 @@ type Path struct { RunOnReady string `json:"runOnReady"` RunOnReadyRestart bool `json:"runOnReadyRestart"` RunOnNotReady string `json:"runOnNotReady"` + RunOnOnline string `json:"runOnOnline"` + RunOnOnlineRestart bool `json:"runOnOnlineRestart"` + RunOnOffline string `json:"runOnOffline"` RunOnRead string `json:"runOnRead"` RunOnReadRestart bool `json:"runOnReadRestart"` RunOnUnread string `json:"runOnUnread"` diff --git a/internal/core/path.go b/internal/core/path.go index 272c5b25..9b835a7d 100644 --- a/internal/core/path.go +++ b/internal/core/path.go @@ -99,6 +99,7 @@ type path struct { onlineTime time.Time onUnDemandHook func(string) onNotReadyHook func() + onOfflineHook func() readers map[defs.Reader]struct{} describeRequestsOnHold []defs.PathDescribeReq readerAddRequestsOnHold []defs.PathAddReaderReq @@ -174,7 +175,7 @@ func (pa *path) isAvailable() bool { } func (pa *path) isOnline() bool { - return pa.source != nil + return pa.onOfflineHook != nil || (pa.source != nil && !pa.conf.AlwaysAvailable) } func (pa *path) run() { @@ -439,7 +440,7 @@ func (pa *path) doSourceStaticSetReady(req defs.PathSourceStaticSetReadyReq) { } if pa.conf.AlwaysAvailable { - pa.onlineTime = time.Now() + pa.setOnline(pa.source.APISourceDescribe(), "") } if pa.conf.HasOnDemandStaticSource() { @@ -457,6 +458,8 @@ func (pa *path) doSourceStaticSetNotReady(req defs.PathSourceStaticSetNotReadyRe if !pa.conf.AlwaysAvailable { pa.setNotAvailable() } else { + pa.setOffline() + err := pa.stream.StartOfflineSubStream() if err != nil { panic("should not happen") @@ -566,7 +569,7 @@ func (pa *path) doAddPublisher(req defs.PathAddPublisherReq) { pa.name) if pa.conf.AlwaysAvailable { - pa.onlineTime = time.Now() + pa.setOnline(req.Author.APISourceDescribe(), req.AccessRequest.Query) } if pa.conf.HasOnDemandPublisher() && pa.onDemandPublisherState != pathOnDemandStateInitial { @@ -822,6 +825,30 @@ func (pa *path) onDemandPublisherStop(reason string) { pa.onDemandPublisherState = pathOnDemandStateInitial } +func (pa *path) setOnline(sourceDesc *defs.APIPathSource, publisherQuery string) { + pa.setOffline() + + pa.onOfflineHook = hooks.OnOnline(hooks.OnOnlineParams{ + Logger: pa, + ExternalCmdPool: pa.externalCmdPool, + Conf: pa.conf, + ExternalCmdEnv: pa.ExternalCmdEnv(), + Desc: sourceDesc, + Query: publisherQuery, + }) + + pa.onlineTime = time.Now() +} + +func (pa *path) setOffline() { + if pa.onOfflineHook == nil { + return + } + + pa.onOfflineHook() + pa.onOfflineHook = nil +} + func (pa *path) setAvailable( source defs.Source, publisherQuery string, @@ -845,10 +872,6 @@ func (pa *path) setAvailable( pa.availableTime = time.Now() - if !pa.conf.AlwaysAvailable { - pa.onlineTime = time.Now() - } - if pa.conf.Record { pa.startRecording() } @@ -867,6 +890,10 @@ func (pa *path) setAvailable( Query: publisherQuery, }) + if !pa.conf.AlwaysAvailable { + pa.setOnline(sourceDesc, publisherQuery) + } + if pa.conf.AlwaysAvailable { pa.Log(logger.Info, "stream is available, %s", defs.MediasInfo(pa.stream.OrigDesc.Medias)) } else { @@ -894,6 +921,7 @@ func (pa *path) consumeOnHoldRequests() { func (pa *path) setNotAvailable() { pa.parent.setPathNotReady(pa) + pa.setOffline() for r := range pa.readers { pa.executeRemoveReader(r) @@ -966,6 +994,8 @@ func (pa *path) executeRemovePublisher() { if !pa.conf.AlwaysAvailable { pa.setNotAvailable() } else { + pa.setOffline() + err := pa.stream.StartOfflineSubStream() if err != nil { panic("should not happen") diff --git a/internal/core/path_test.go b/internal/core/path_test.go index 5b83c2c4..64c4632e 100644 --- a/internal/core/path_test.go +++ b/internal/core/path_test.go @@ -366,6 +366,82 @@ func TestPathRunOnReadyQueryInjection(t *testing.T) { } } +func TestPathRunOnOnline(t *testing.T) { + onOnline := filepath.Join(t.TempDir(), "on_online") + onOffline := filepath.Join(t.TempDir(), "on_offline") + + func() { + p, ok := newInstance(t, fmt.Sprintf("rtmp: no\n"+ + "hls: no\n"+ + "webrtc: no\n"+ + "paths:\n"+ + " test:\n"+ + " alwaysAvailable: yes\n"+ + " alwaysAvailableTracks:\n"+ + " - codec: H264\n"+ + " runOnOnline: sh -c 'echo \"$MTX_PATH $MTX_QUERY $MTX_SOURCE_TYPE $MTX_SOURCE_ID $RTSP_PORT $G1\" "+ + "> %s; while true; do sleep 1; done'\n"+ + " runOnOffline: sh -c 'echo \"$MTX_PATH $MTX_QUERY $MTX_SOURCE_TYPE $MTX_SOURCE_ID $RTSP_PORT $G1\" "+ + "> %s'\n", + onOnline, onOffline)) + require.Equal(t, true, ok) + defer p.Close() + + _, err := os.Stat(onOnline) + require.ErrorIs(t, err, os.ErrNotExist) + + c := gortsplib.Client{} + + err = c.StartRecording( + "rtsp://localhost:8554/test?query=value", + &description.Session{Medias: []*description.Media{test.UniqueMediaH264()}}) + require.NoError(t, err) + + for { + _, err = os.Stat(onOnline) + if err == nil { + break + } + require.ErrorIs(t, err, os.ErrNotExist) + time.Sleep(50 * time.Millisecond) + } + + _, err = os.Stat(onOffline) + require.ErrorIs(t, err, os.ErrNotExist) + + c.Close() + + for { + _, err = os.Stat(onOffline) + if err == nil { + break + } + require.ErrorIs(t, err, os.ErrNotExist) + time.Sleep(50 * time.Millisecond) + } + }() + + byts, err := os.ReadFile(onOnline) + require.NoError(t, err) + fields := strings.Split(string(byts[:len(byts)-1]), " ") + require.Equal(t, "test", fields[0]) + require.Equal(t, "query%3Dvalue", fields[1]) + require.Equal(t, "rtspSession", fields[2]) + require.NotEmpty(t, fields[3]) + require.Equal(t, "8554", fields[4]) + require.Equal(t, "", fields[5]) + + byts, err = os.ReadFile(onOffline) + require.NoError(t, err) + fields = strings.Split(string(byts[:len(byts)-1]), " ") + require.Equal(t, "test", fields[0]) + require.Equal(t, "query%3Dvalue", fields[1]) + require.Equal(t, "rtspSession", fields[2]) + require.NotEmpty(t, fields[3]) + require.Equal(t, "8554", fields[4]) + require.Equal(t, "", fields[5]) +} + func TestPathRunOnRead(t *testing.T) { serverCertFpath := test.CreateTempFile(t, test.TLSCertPub) serverKeyFpath := test.CreateTempFile(t, test.TLSCertKey) diff --git a/internal/hooks/on_online.go b/internal/hooks/on_online.go new file mode 100644 index 00000000..580e8a5c --- /dev/null +++ b/internal/hooks/on_online.go @@ -0,0 +1,69 @@ +package hooks + +import ( + "net/url" + + "github.com/bluenviron/mediamtx/internal/conf" + "github.com/bluenviron/mediamtx/internal/defs" + "github.com/bluenviron/mediamtx/internal/externalcmd" + "github.com/bluenviron/mediamtx/internal/logger" +) + +// OnOnlineParams are the parameters of OnOnline. +type OnOnlineParams struct { + Logger logger.Writer + ExternalCmdPool *externalcmd.Pool + Conf *conf.Path + ExternalCmdEnv externalcmd.Environment + Desc *defs.APIPathSource + Query string +} + +// OnOnline is the OnOnline hook. +func OnOnline(params OnOnlineParams) func() { + var env externalcmd.Environment + var onOnlineCmd *externalcmd.Cmd + + if params.Conf.RunOnOnline != "" || params.Conf.RunOnOffline != "" { + env = params.ExternalCmdEnv + env["MTX_QUERY"] = url.QueryEscape(params.Query) + if params.Desc != nil { + env["MTX_SOURCE_TYPE"] = string(params.Desc.Type) + env["MTX_SOURCE_ID"] = params.Desc.ID + } + } + + if params.Conf.RunOnOnline != "" { + params.Logger.Log(logger.Info, "runOnOnline command started") + onOnlineCmd = &externalcmd.Cmd{ + Pool: params.ExternalCmdPool, + Cmdstr: params.Conf.RunOnOnline, + Restart: params.Conf.RunOnOnlineRestart, + Env: env, + OnExit: func(err error) { + params.Logger.Log(logger.Info, "runOnOnline command exited: %v", err) + }, + } + onOnlineCmd.Start() + } + + return func() { + if onOnlineCmd != nil { + onlineCmd := onOnlineCmd + onOnlineCmd = nil + onlineCmd.Close() + params.Logger.Log(logger.Info, "runOnOnline command stopped") + } + + if params.Conf.RunOnOffline != "" { + params.Logger.Log(logger.Info, "runOnOffline command launched") + cmd := &externalcmd.Cmd{ + Pool: params.ExternalCmdPool, + Cmdstr: params.Conf.RunOnOffline, + Restart: false, + Env: env, + } + cmd.Start() + } + } +} diff --git a/mediamtx.yml b/mediamtx.yml index 78b3faba..0b871ff6 100644 --- a/mediamtx.yml +++ b/mediamtx.yml @@ -758,8 +758,7 @@ pathDefaults: # Environment variables are the same as runOnDemand. runOnUnDemand: - # Command to run when the stream is ready to be read, whenever it is - # published by a client or pulled from a server / camera. + # Command to run when the stream is available to be read. # This is terminated with SIGINT when the stream is not ready anymore. # The following environment variables are available: # * MTX_PATH: path name @@ -776,6 +775,24 @@ pathDefaults: # Environment variables are the same as runOnReady. runOnNotReady: + # Command to run when the stream is online, which means + # that the stream is available and provided by a online source (not an offline segment). + # This is terminated with SIGINT when the stream is not online anymore. + # The following environment variables are available: + # * MTX_PATH: path name + # * MTX_QUERY: query parameters (passed by publisher) (url-encoded) + # * MTX_SOURCE_TYPE: source type + # * MTX_SOURCE_ID: source ID + # * RTSP_PORT: RTSP server port + # * G1, G2, ...: regular expression groups, if path name is + # a regular expression. + runOnOnline: + # Restart the command if it exits. + runOnOnlineRestart: false + # Command to run when the stream is not online anymore. + # Environment variables are the same as runOnOnline. + runOnOffline: + # Command to run when a client starts reading. # This is terminated with SIGINT when a client stops reading. # The following environment variables are available: