From d6dba0ef80d42c68698a20f415775ba105537b28 Mon Sep 17 00:00:00 2001 From: Eli Streams Date: Thu, 13 Mar 2025 03:41:52 +0000 Subject: [PATCH] gstreamer: reduce memory leakage & clean up See merge request streamplace/streamplace!107 Changelog: feature --- .gitlab-ci.yml | 3 +- go.mod | 12 +- go.sum | 20 +-- pkg/api/playback.go | 12 ++ pkg/cmd/streamplace.go | 9 ++ pkg/cmd/whip.go | 72 ++++----- pkg/media/bus_handler.go | 42 ++++++ pkg/media/concat.go | 2 + pkg/media/gstreamer.go | 292 ++++++++++++------------------------- pkg/media/media.go | 7 +- pkg/media/webrtc.go | 46 +----- pkg/thumbnail/thumbnail.go | 24 +++ test/leak-check.sh | 6 +- 13 files changed, 252 insertions(+), 295 deletions(-) create mode 100644 pkg/media/bus_handler.go create mode 100644 pkg/thumbnail/thumbnail.go diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index bb79915ec..aee981037 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -155,7 +155,8 @@ build-mac: - 'echo "xcodeproj --version: $(xcodeproj --version)"' - "$(which xcodeproj) --version" - > - brew install python@3.11 + export GOTOOLCHAIN=go1.23.3 + && brew install python@3.11 && python3.11 -m pip install virtualenv && python3.11 -m virtualenv ~/venv && source ~/venv/bin/activate diff --git a/go.mod b/go.mod index 1208cac15..a4fcbf5f3 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module stream.place/streamplace -go 1.23 +go 1.23.1 toolchain go1.23.3 @@ -21,8 +21,8 @@ require ( github.com/dunglas/httpsfv v1.0.2 github.com/ethereum/go-ethereum v1.14.7 github.com/go-git/go-git/v5 v5.12.0 - github.com/go-gst/go-glib v1.3.0 - github.com/go-gst/go-gst v1.3.0 + github.com/go-gst/go-glib v1.4.0 + github.com/go-gst/go-gst v1.4.0 github.com/golang/freetype v0.0.0-20170609003504-e2365dfdc4a0 github.com/golang/glog v1.2.0 github.com/google/uuid v1.6.0 @@ -48,7 +48,7 @@ require ( github.com/stretchr/testify v1.10.0 github.com/whyrusleeping/cbor-gen v0.2.1-0.20241030202151-b7a6831be65e gitlab.com/gitlab-org/release-cli v0.18.0 - golang.org/x/exp v0.0.0-20240719175910-8a7402abbf56 + golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 golang.org/x/image v0.22.0 golang.org/x/net v0.33.0 golang.org/x/sync v0.10.0 @@ -236,11 +236,11 @@ require ( go.uber.org/multierr v1.11.0 // indirect go.uber.org/zap v1.26.0 // indirect golang.org/x/crypto v0.31.0 // indirect - golang.org/x/mod v0.20.0 // indirect + golang.org/x/mod v0.21.0 // indirect golang.org/x/oauth2 v0.21.0 // indirect golang.org/x/text v0.21.0 // indirect golang.org/x/time v0.5.0 // indirect - golang.org/x/tools v0.24.0 // indirect + golang.org/x/tools v0.25.0 // indirect google.golang.org/appengine/v2 v2.0.2 // indirect google.golang.org/genproto v0.0.0-20240722135656-d784300faade // indirect google.golang.org/genproto/googleapis/api v0.0.0-20240701130421-f6361c86f094 // indirect diff --git a/go.sum b/go.sum index 8858402c8..2173084f1 100644 --- a/go.sum +++ b/go.sum @@ -175,10 +175,10 @@ github.com/go-git/go-git-fixtures/v4 v4.3.2-0.20231010084843-55a94097c399 h1:eMj github.com/go-git/go-git-fixtures/v4 v4.3.2-0.20231010084843-55a94097c399/go.mod h1:1OCfN199q1Jm3HZlxleg+Dw/mwps2Wbk9frAWm+4FII= github.com/go-git/go-git/v5 v5.12.0 h1:7Md+ndsjrzZxbddRDZjF14qK+NN56sy6wkqaVrjZtys= github.com/go-git/go-git/v5 v5.12.0/go.mod h1:FTM9VKtnI2m65hNI/TenDDDnUf2Q9FHnXYjuz9i5OEY= -github.com/go-gst/go-glib v1.3.0 h1:u+mPUdLmrDFA/MskIxInJY+M0O1RSkHeZYggnJGWlPk= -github.com/go-gst/go-glib v1.3.0/go.mod h1:JybIYeoHNwCkHGaBf1fHNIaM4sQTrJPkPLsi7dmPNOU= -github.com/go-gst/go-gst v1.3.0 h1:z4mQ7CNJXd6ZfkibzIT9kZKwtgEFJo7jJGlX9cXFzz0= -github.com/go-gst/go-gst v1.3.0/go.mod h1:2li6ghiCBz7/R6DA7itVto3gsYh0QKicwSxEefNVYqE= +github.com/go-gst/go-glib v1.4.0 h1:FB2uVfB0uqz7/M6EaDdWWlBZRQpvFAbWfL7drdw8lAE= +github.com/go-gst/go-glib v1.4.0/go.mod h1:GUIpWmkxQ1/eL+FYSjKpLDyTZx6Vgd9nNXt8dA31d5M= +github.com/go-gst/go-gst v1.4.0 h1:EikB43u4c3wc8d2RzlFRSfIGIXYzDy6Zls2vJqrG2BU= +github.com/go-gst/go-gst v1.4.0/go.mod h1:p8TLGtOxJLcrp6PCkTPdnanwWBxPZvYiHDbuSuwgO3c= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= @@ -678,8 +678,8 @@ golang.org/x/crypto v0.7.0/go.mod h1:pYwdfH91IfpZVANVyUOhSIPZaFoJGxTFbZhFTx+dXZU golang.org/x/crypto v0.31.0 h1:ihbySMvVjLAeSH1IbfcRTkD/iNscyz8rGzjF/E5hV6U= golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk= golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= -golang.org/x/exp v0.0.0-20240719175910-8a7402abbf56 h1:2dVuKD2vS7b0QIHQbpyTISPd0LeHDbnYEryqj5Q1ug8= -golang.org/x/exp v0.0.0-20240719175910-8a7402abbf56/go.mod h1:M4RDyNAINzryxdtnbRXRL/OHtkFuWGRjvuhBJpk2IlY= +golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 h1:e66Fs6Z+fZTbFBAxKfP3PALWBtpfqks2bwGcexMxgtk= +golang.org/x/exp v0.0.0-20240909161429-701f63a606c0/go.mod h1:2TbTHSBQa924w8M6Xs1QcRcFwyucIwBGpK1p2f1YFFY= golang.org/x/image v0.22.0 h1:UtK5yLUzilVrkjMAZAZ34DXGpASN8i8pj8g+O+yd10g= golang.org/x/image v0.22.0/go.mod h1:9hPFhljd4zZ1GNSIZJ49sqbp45GKK9t6w+iXvGqZUz4= golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE= @@ -692,8 +692,8 @@ golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.4.2/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= -golang.org/x/mod v0.20.0 h1:utOm6MM3R3dnawAiJgn0y+xvuYRsm1RKM/4giyfDgV0= -golang.org/x/mod v0.20.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= +golang.org/x/mod v0.21.0 h1:vvrHzRwRfVKSiLrG+d4FMl/Qi4ukBCE6kZlTUkDYRT0= +golang.org/x/mod v0.21.0/go.mod h1:6SkKJ3Xj0I0BrPOZoBy3bdMptDDU9oJrpohJ3eWZ1fY= golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20190213061140-3a22650c66bd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= @@ -782,8 +782,8 @@ golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4f golang.org/x/tools v0.1.5/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= -golang.org/x/tools v0.24.0 h1:J1shsA93PJUEVaUSaay7UXAyE8aimq3GW0pjlolpa24= -golang.org/x/tools v0.24.0/go.mod h1:YhNqVBIfWHdzvTLs0d8LCuMhkKUgSUKldakyV7W/WDQ= +golang.org/x/tools v0.25.0 h1:oFU9pkj/iJgs+0DT+VMHrx+oBKs/LJMV+Uvg78sl+fE= +golang.org/x/tools v0.25.0/go.mod h1:/vtpO8WL1N9cQC3FN5zPqb//fRXskFHbLKk4OW1Q7rg= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= diff --git a/pkg/api/playback.go b/pkg/api/playback.go index f6fab9ddb..0fc687398 100644 --- a/pkg/api/playback.go +++ b/pkg/api/playback.go @@ -81,6 +81,12 @@ func (a *StreamplaceAPI) HandleMP4Playback(ctx context.Context) httprouter.Handl g.Go(func() error { return a.MediaManager.SegmentToMP4(ctx, user, bufw) }) + g.Go(func() error { + <-ctx.Done() + pr.Close() + pw.Close() + return nil + }) g.Go(func() error { time.Sleep(time.Duration(delayMS) * time.Millisecond) _, err := io.Copy(w, pr) @@ -126,6 +132,12 @@ func (a *StreamplaceAPI) HandleMKVPlayback(ctx context.Context) httprouter.Handl g.Go(func() error { return a.MediaManager.SegmentToMKV(ctx, user, bufw) }) + g.Go(func() error { + <-ctx.Done() + pr.Close() + pw.Close() + return nil + }) g.Go(func() error { time.Sleep(time.Duration(delayMS) * time.Millisecond) _, err := io.Copy(w, pr) diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index a63ea5966..830fda4e1 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -28,6 +28,7 @@ import ( "stream.place/streamplace/pkg/replication/boring" v0 "stream.place/streamplace/pkg/schema/v0" "stream.place/streamplace/pkg/spmetrics" + "stream.place/streamplace/pkg/thumbnail" "github.com/ThalesGroup/crypto11" _ "github.com/go-gst/go-glib/glib" @@ -334,6 +335,13 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { } go func() { err := func() error { + lock := thumbnail.GetThumbnailLock(not.Segment.RepoDID) + locked := lock.TryLock() + if !locked { + // we're already generating a thumbnail for this user, skip + return nil + } + defer lock.Unlock() oldThumb, err := mod.LatestThumbnailForUser(not.Segment.RepoDID) if err != nil { return err @@ -348,6 +356,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { if err != nil { return err } + defer fd.Close() err = mm.Thumbnail(ctx, r, fd) if err != nil { return err diff --git a/pkg/cmd/whip.go b/pkg/cmd/whip.go index dbb29e9d3..f130ae2bb 100644 --- a/pkg/cmd/whip.go +++ b/pkg/cmd/whip.go @@ -14,9 +14,10 @@ import ( "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" "github.com/pion/webrtc/v4" - "github.com/pion/webrtc/v4/pkg/media" + pionmedia "github.com/pion/webrtc/v4/pkg/media" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/media" ) func WHIP() error { @@ -26,6 +27,7 @@ func WHIP() error { duration := fs.Duration("duration", 0, "duration of the stream") file := fs.String("file", "", "file to stream (needs to be an MP4 containing H264 video and Opus audio)") endpoint := fs.String("endpoint", "http://127.0.0.1:38080", "endpoint to send the WHIP request to") + freezeAfter := fs.Duration("freeze-after", 0, "freeze the stream after the given duration") err := fs.Parse(os.Args[2:]) if *file == "" { return fmt.Errorf("file is required") @@ -42,20 +44,22 @@ func WHIP() error { } w := &WHIPClient{ - StreamKey: *streamKey, - File: *file, - Endpoint: *endpoint, - Count: *count, + StreamKey: *streamKey, + File: *file, + Endpoint: *endpoint, + Count: *count, + FreezeAfter: *freezeAfter, } return w.WHIP(ctx) } type WHIPClient struct { - StreamKey string - File string - Endpoint string - Count int + StreamKey string + File string + Endpoint string + Count int + FreezeAfter time.Duration } var failureStates = []webrtc.ICEConnectionState{ @@ -316,18 +320,20 @@ func (w *WHIPClient) WHIP(ctx context.Context) error { accumulators[i] += duration - for _, conn := range conns { - if trackType == "video" { - if err := conn.videoTrack.WriteSample(media.Sample{Data: samples, Duration: duration}); err != nil { - log.Log(ctx, "error writing video sample", "error", err) - errCh <- err - return gst.FlowError - } - } else { - if err := conn.audioTrack.WriteSample(media.Sample{Data: samples, Duration: duration}); err != nil { - log.Log(ctx, "error writing video sample", "error", err) - errCh <- err - return gst.FlowError + if w.FreezeAfter == 0 || time.Since(startTime) < w.FreezeAfter { + for _, conn := range conns { + if trackType == "video" { + if err := conn.videoTrack.WriteSample(pionmedia.Sample{Data: samples, Duration: duration}); err != nil { + log.Log(ctx, "error writing video sample", "error", err) + errCh <- err + return gst.FlowError + } + } else { + if err := conn.audioTrack.WriteSample(pionmedia.Sample{Data: samples, Duration: duration}); err != nil { + log.Log(ctx, "error writing video sample", "error", err) + errCh <- err + return gst.FlowError + } } } } @@ -338,26 +344,10 @@ func (w *WHIPClient) WHIP(ctx context.Context) error { }(i) } - ok := pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - log.Log(ctx, "got gst.MessageEOS, exiting") - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) - if !ok { - return fmt.Errorf("failed to add watch to pipeline bus") - } + go func() { + media.HandleBusMessages(ctx, pipeline) + cancel() + }() if err = pipeline.SetState(gst.StatePlaying); err != nil { return err diff --git a/pkg/media/bus_handler.go b/pkg/media/bus_handler.go new file mode 100644 index 000000000..7e8ba03fe --- /dev/null +++ b/pkg/media/bus_handler.go @@ -0,0 +1,42 @@ +package media + +import ( + "context" + "time" + + "github.com/go-gst/go-gst/gst" + "stream.place/streamplace/pkg/log" +) + +func HandleBusMessages(ctx context.Context, pipeline *gst.Pipeline) { + HandleBusMessagesCustom(ctx, pipeline, nil) +} + +func HandleBusMessagesCustom(ctx context.Context, pipeline *gst.Pipeline, handler func(msg *gst.Message)) { + for { + if ctx.Err() != nil { + return + } + msg := pipeline.GetPipelineBus().PopMessage(gst.ClockTime(time.Second * 1)) + if msg == nil { + continue + } + if handler != nil { + handler(msg) + } + switch msg.Type() { + case gst.MessageEOS: // When end-of-stream is received flush the pipeline and stop the main loop + log.Log(ctx, "got gst.MessageEOS, exiting") + return + case gst.MessageError: // Error messages are always fatal + err := msg.ParseError() + log.Error(ctx, "gstreamer error", "error", err.Error()) + if debug := err.DebugString(); debug != "" { + log.Debug(ctx, "gstreamer debug", "message", debug) + } + return + default: + log.Debug(ctx, msg.String()) + } + } +} diff --git a/pkg/media/concat.go b/pkg/media/concat.go index 7b06f00b5..3c5928be0 100644 --- a/pkg/media/concat.go +++ b/pkg/media/concat.go @@ -153,6 +153,8 @@ func ConcatStream(ctx context.Context, pipeline *gst.Pipeline, user string, stre go func() { select { case <-ctx.Done(): + pr.Close() + pw.Close() return case fullpath := <-allFiles: if fullpath == "" { diff --git a/pkg/media/gstreamer.go b/pkg/media/gstreamer.go index 45b0fd3e6..79b15a819 100644 --- a/pkg/media/gstreamer.go +++ b/pkg/media/gstreamer.go @@ -39,6 +39,9 @@ func SafePipe() (*os.File, *os.File, func(), error) { func ReaderNeedData(ctx context.Context, input io.Reader) func(self *app.Source, length uint) { return func(self *app.Source, length uint) { + if ctx.Err() != nil { + return + } bs := make([]byte, length) read, err := input.Read(bs) if err != nil { @@ -70,7 +73,6 @@ func WriterNewSample(ctx context.Context, output io.Writer) func(sink *app.Sink) if sample == nil { return gst.FlowOK } - // defer sample.Unref() // Retrieve the buffer from the sample. buffer := sample.GetBuffer() @@ -86,9 +88,6 @@ func WriterNewSample(ctx context.Context, output io.Writer) func(sink *app.Sink) } func AddOpusToMKV(ctx context.Context, input io.Reader, output io.Writer) error { - - mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) - pipelineSlice := []string{ "appsrc name=appsrc ! matroskademux name=demux", "matroskamux name=mux ! appsink name=appsink", @@ -130,36 +129,16 @@ func AddOpusToMKV(ctx context.Context, input io.Reader, output io.Writer) error }) go func() { - <-ctx.Done() - pipeline.BlockSetState(gst.StateNull) - mainLoop.Quit() + HandleBusMessages(ctx, pipeline) + cancel() }() - // Add a message handler to the pipeline bus, printing interesting information to the console. - pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - log.Debug(ctx, "got EOS") - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) - // Start the pipeline pipeline.SetState(gst.StatePlaying) - mainLoop.Run() - log.Log(ctx, "main loop complete") + <-ctx.Done() + + pipeline.BlockSetState(gst.StateNull) return nil } @@ -269,7 +248,6 @@ func SelfTest(ctx context.Context) error { // #EXT-X-ENDLIST func (mm *MediaManager) ToHLS(ctx context.Context, input io.Reader, m3u8 *M3U8) error { - mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) ctx = log.WithLogValues(ctx, "GStreamerFunc", "ToHLS") splitmuxsink, err := gst.NewElementWithProperties("splitmuxsink", map[string]any{ @@ -430,68 +408,55 @@ func (mm *MediaManager) ToHLS(ctx context.Context, input io.Reader, m3u8 *M3U8) ctx, cancel := context.WithCancel(ctx) defer cancel() go func() { - <-ctx.Done() - pipeline.BlockSetState(gst.StateNull) - mainLoop.Quit() - }() - - pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Debug(ctx, "gstreamer debug", "message", debug) - } - cancel() - case gst.MessageElement: - structure := msg.GetStructure() - name := structure.Name() - if name == "splitmuxsink-fragment-opened" { - runningTime, err := structure.GetValue("running-time") - if err != nil { - log.Warn(ctx, "splitmuxsink-fragment-opened error", "error", err) - return true - } - runningTimeInt, ok := runningTime.(uint64) - if !ok { - log.Warn(ctx, "splitmuxsink-fragment-opened not a uint64") - return true - } - m3u8.FragmentOpened(ctx, runningTimeInt) - } - if name == "splitmuxsink-fragment-closed" { - runningTime, err := structure.GetValue("running-time") - if err != nil { - log.Warn(ctx, "splitmuxsink-fragment-closed error", "error", err) - return true + HandleBusMessagesCustom(ctx, pipeline, func(msg *gst.Message) { + switch msg.Type() { + case gst.MessageElement: + structure := msg.GetStructure() + name := structure.Name() + if name == "splitmuxsink-fragment-opened" { + runningTime, err := structure.GetValue("running-time") + if err != nil { + log.Warn(ctx, "splitmuxsink-fragment-opened error", "error", err) + cancel() + } + runningTimeInt, ok := runningTime.(uint64) + if !ok { + log.Warn(ctx, "splitmuxsink-fragment-opened not a uint64") + cancel() + } + m3u8.FragmentOpened(ctx, runningTimeInt) } - runningTimeInt, ok := runningTime.(uint64) - if !ok { - log.Warn(ctx, "splitmuxsink-fragment-closed not a uint64") - return true + if name == "splitmuxsink-fragment-closed" { + runningTime, err := structure.GetValue("running-time") + if err != nil { + log.Warn(ctx, "splitmuxsink-fragment-closed error", "error", err) + cancel() + } + runningTimeInt, ok := runningTime.(uint64) + if !ok { + log.Warn(ctx, "splitmuxsink-fragment-closed not a uint64") + cancel() + } + m3u8.FragmentClosed(ctx, runningTimeInt) } - m3u8.FragmentClosed(ctx, runningTimeInt) } - default: - log.Debug(ctx, msg.String()) - } - return true - }) + }) + cancel() + }() // Start the pipeline pipeline.SetState(gst.StatePlaying) - mainLoop.Run() - log.Log(ctx, "main loop complete") + <-ctx.Done() + + pipeline.BlockSetState(gst.StateNull) return nil } func (mm *MediaManager) IngestStream(ctx context.Context, input io.Reader, ms MediaSigner) error { + ctx, cancel := context.WithCancel(ctx) + defer cancel() pipelineSlice := []string{ "appsrc name=streamsrc ! matroskademux name=demux", "demux. ! queue ! h264parse name=parse", @@ -538,32 +503,17 @@ func (mm *MediaManager) IngestStream(ctx context.Context, input io.Reader, ms Me return err } - mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) - - pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - mainLoop.Quit() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - mainLoop.Quit() - default: - log.Debug(ctx, msg.String()) - } - return true - }) + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() err = pipeline.SetState(gst.StatePlaying) if err != nil { return err } - mainLoop.Run() + <-ctx.Done() return nil } @@ -674,23 +624,10 @@ func (mm *MediaManager) TestSource(ctx context.Context, ms MediaSigner) error { mainLoop.Quit() }() - pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Log(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() // Start the pipeline pipeline.SetState(gst.StatePlaying) @@ -729,6 +666,23 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) return nil, fmt.Errorf("failed to get audio pad") } + resetTimer := make(chan struct{}) + + go func() { + for { + select { + case <-ctx.Done(): + return + case <-resetTimer: + continue + case <-time.After(time.Second * 10): + log.Warn(ctx, "no new segment for 10 seconds") + elem.ErrorMessage(gst.DomainCore, gst.CoreErrorFailed, "No new segment for 10 seconds", "No new segment for 10 seconds (debug)") + return + } + } + }() + elem.Connect("sink-added", func(split, sinkEle *gst.Element) { buf := &bytes.Buffer{} appsink := app.SinkFromElement(sinkEle) @@ -738,6 +692,7 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) appsink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: WriterNewSample(ctx, buf), EOSFunc: func(sink *app.Sink) { + resetTimer <- struct{}{} bs, err := ms.SignMP4(ctx, bytes.NewReader(buf.Bytes()), time.Now().UnixMilli()) if err != nil { log.Error(ctx, "error signing segment", "error", err) @@ -757,15 +712,16 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) func (mm *MediaManager) Thumbnail(ctx context.Context, r io.Reader, w io.Writer) error { ctx = log.WithLogValues(ctx, "function", "Thumbnail") - mainLoop := glib.NewMainLoop(glib.MainContextDefault(), false) + ctx, cancel := context.WithCancel(ctx) + defer cancel() pipelineSlice := []string{ - "appsrc name=appsrc ! qtdemux ! decodebin ! videoconvert ! videoscale ! video/x-raw,width=[1,200],height=[1,200],pixel-aspect-ratio=1/1 ! pngenc ! appsink name=appsink", + "appsrc name=appsrc ! qtdemux ! decodebin ! videoconvert ! videoscale ! video/x-raw,width=[1,720],height=[1,720],pixel-aspect-ratio=1/1 ! pngenc snapshot=true ! appsink name=appsink", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) if err != nil { - return fmt.Errorf("error creating TestSource pipeline: %w", err) + return fmt.Errorf("error creating Thumbnail pipeline: %w", err) } appsrc, err := pipeline.GetElementByName("appsrc") if err != nil { @@ -777,31 +733,15 @@ func (mm *MediaManager) Thumbnail(ctx context.Context, r io.Reader, w io.Writer) NeedDataFunc: ReaderNeedData(ctx, r), }) - ctx, cancel := context.WithCancel(ctx) - defer cancel() - appsink, err := pipeline.GetElementByName("appsink") if err != nil { return err } - pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Log(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() sink := app.SinkFromElement(appsink) sink.SetCallbacks(&app.SinkCallbacks{ @@ -813,13 +753,9 @@ func (mm *MediaManager) Thumbnail(ctx context.Context, r io.Reader, w io.Writer) pipeline.SetState(gst.StatePlaying) - go func() { - <-ctx.Done() - pipeline.BlockSetState(gst.StateNull) - mainLoop.Quit() - }() + <-ctx.Done() - mainLoop.Run() + pipeline.BlockSetState(gst.StateNull) return nil } @@ -845,26 +781,10 @@ func (mm *MediaManager) MP4Playback(ctx context.Context, user string, w io.Write return fmt.Errorf("failed to create GStreamer pipeline: %w", err) } - ok := pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - log.Log(ctx, "got gst.MessageEOS, exiting") - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) - if !ok { - return fmt.Errorf("failed to add watch to pipeline bus") - } + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() outputQueue, done, err := ConcatStream(ctx, pipeline, user, mm) if err != nil { @@ -897,11 +817,6 @@ func (mm *MediaManager) MP4Playback(ctx context.Context, user string, w io.Write return fmt.Errorf("failed to link output queue to audio parse: %w", err) } - go func() { - <-ctx.Done() - pipeline.BlockSetState(gst.StateNull) - }() - go func() { ticker := time.NewTicker(time.Second * 1) for { @@ -931,6 +846,9 @@ func (mm *MediaManager) MP4Playback(ctx context.Context, user string, w io.Write pipeline.SetState(gst.StatePlaying) <-ctx.Done() + + pipeline.BlockSetState(gst.StateNull) + return nil } @@ -955,26 +873,10 @@ func (mm *MediaManager) MKVPlayback(ctx context.Context, user string, w io.Write return fmt.Errorf("failed to create GStreamer pipeline: %w", err) } - ok := pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - log.Log(ctx, "got gst.MessageEOS, exiting") - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) - if !ok { - return fmt.Errorf("failed to add watch to pipeline bus") - } + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() outputQueue, done, err := ConcatStream(ctx, pipeline, user, mm) if err != nil { @@ -1007,11 +909,6 @@ func (mm *MediaManager) MKVPlayback(ctx context.Context, user string, w io.Write return fmt.Errorf("failed to link output queue to audio parse: %w", err) } - go func() { - <-ctx.Done() - pipeline.BlockSetState(gst.StateNull) - }() - go func() { ticker := time.NewTicker(time.Second * 1) for { @@ -1041,5 +938,8 @@ func (mm *MediaManager) MKVPlayback(ctx context.Context, user string, w io.Write pipeline.SetState(gst.StatePlaying) <-ctx.Done() + + pipeline.BlockSetState(gst.StateNull) + return nil } diff --git a/pkg/media/media.go b/pkg/media/media.go index cbb77d8c3..2c87c41ab 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -76,7 +76,7 @@ func MakeMediaManager(ctx context.Context, cli *config.CLI, signer crypto.Signer } // replacement for os.Pipe that works on windows -func (mm *MediaManager) HTTPPipe() (string, io.Reader, func(), error) { +func (mm *MediaManager) HTTPPipe() (string, io.ReadCloser, func(), error) { uu, err := uuid.NewV7() if err != nil { return "", nil, nil, err @@ -249,6 +249,11 @@ func (mm *MediaManager) SegmentToStream(ctx context.Context, user string, muxer }, } g, _ := errgroup.WithContext(ctx) + g.Go(func() error { + <-ctx.Done() + or.Close() + return nil + }) g.Go(func() error { _, err := tc.Transcode(in, out) tc.StopTranscoder() diff --git a/pkg/media/webrtc.go b/pkg/media/webrtc.go index 121510262..97546afb8 100644 --- a/pkg/media/webrtc.go +++ b/pkg/media/webrtc.go @@ -44,26 +44,10 @@ func (mm *MediaManager) WebRTCPlayback(ctx context.Context, user string, offer * return nil, fmt.Errorf("failed to create GStreamer pipeline: %w", err) } - ok := pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - log.Log(ctx, "got gst.MessageEOS, exiting") - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Log(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) - if !ok { - return nil, fmt.Errorf("failed to add watch to pipeline bus") - } + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() outputQueue, done, err := ConcatStream(ctx, pipeline, user, mm) if err != nil { @@ -415,24 +399,10 @@ func (mm *MediaManager) WebRTCIngest(ctx context.Context, offer *webrtc.SessionD return nil, fmt.Errorf("failed to create GStreamer pipeline: %w", err) } - pipeline.GetPipelineBus().AddWatch(func(msg *gst.Message) bool { - switch msg.Type() { - - case gst.MessageEOS: // When end-of-stream is received flush the pipeling and stop the main loop - log.Debug(ctx, "got gst.MessageEOS, exiting") - cancel() - case gst.MessageError: // Error messages are always fatal - err := msg.ParseError() - log.Error(ctx, "gstreamer error", "error", err.Error()) - if debug := err.DebugString(); debug != "" { - log.Debug(ctx, "gstreamer debug", "message", debug) - } - cancel() - default: - log.Debug(ctx, msg.String()) - } - return true - }) + go func() { + HandleBusMessages(ctx, pipeline) + cancel() + }() queue, err := pipeline.GetElementByName("queue") if err != nil { diff --git a/pkg/thumbnail/thumbnail.go b/pkg/thumbnail/thumbnail.go new file mode 100644 index 000000000..ceed89c9f --- /dev/null +++ b/pkg/thumbnail/thumbnail.go @@ -0,0 +1,24 @@ +package thumbnail + +import "sync" + +var thumbnailLocks = struct { + sync.Mutex + locks map[string]*sync.Mutex +}{ + locks: make(map[string]*sync.Mutex), +} + +// GetThumbnailLock returns a mutex for the given user +func GetThumbnailLock(handle string) *sync.Mutex { + thumbnailLocks.Lock() + defer thumbnailLocks.Unlock() + + if lock, exists := thumbnailLocks.locks[handle]; exists { + return lock + } + + lock := &sync.Mutex{} + thumbnailLocks.locks[handle] = lock + return lock +} diff --git a/test/leak-check.sh b/test/leak-check.sh index 274f23a48..c65f58e25 100755 --- a/test/leak-check.sh +++ b/test/leak-check.sh @@ -8,12 +8,14 @@ DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" && cd .. && pwd )" TMPDIR="$(mktemp -d)" cd $TMPDIR -MALLOC_CONF=prof_leak:true,lg_prof_sample:0,prof_final:true LD_PRELOAD=/usr/lib/x86_64-linux-gnu/libjemalloc.so.2 "$DIR/build-linux-amd64/streamplace" --no-firehose --wide-open & +MALLOC_CONF=prof_leak:true,lg_prof_sample:0,prof_final:true LD_PRELOAD=/usr/lib/x86_64-linux-gnu/libjemalloc.so.2 \ + "$DIR/build-linux-amd64/streamplace" --no-firehose --wide-open & STREAMPLACE_PID=$! sleep 3 -"$DIR/build-linux-amd64/streamplace" whip --count=3 --duration=30s || true +"$DIR/build-linux-amd64/streamplace" whip --count=3 --duration=90s \ + --file=$HOME/testvids/RocketLeague_1h55m_1sGOP_1080p60_NoBframes.mp4 || true sleep 5 curl -X POST http://127.0.0.1:39090/gc -- 2.51.2