diff --git a/pkg/media/bus_handler.go b/pkg/media/bus_handler.go index 87eed3212..54d4f2dcf 100644 --- a/pkg/media/bus_handler.go +++ b/pkg/media/bus_handler.go @@ -2,6 +2,7 @@ package media import ( "context" + "errors" "fmt" "time" @@ -9,6 +10,12 @@ import ( "stream.place/streamplace/pkg/log" ) +// ErrPipelineDone is a sentinel a pipeline element can raise (via +// pipeline.Error) to signal it finished its job and the bus handler should +// exit cleanly rather than treat it as a failure — e.g. the thumbnailer once +// it has grabbed its frame. +var ErrPipelineDone = errors.New("pipeline done") + func HandleBusMessages(ctx context.Context, pipeline *gst.Pipeline) error { return HandleBusMessagesCustom(ctx, pipeline, nil) } @@ -31,8 +38,8 @@ func HandleBusMessagesCustom(ctx context.Context, pipeline *gst.Pipeline, handle return nil case gst.MessageError: // Error messages are always fatal err := msg.ParseError() - if err.Error() == fmt.Sprintf("%s: %s", ErrConcatDone.Error(), ErrConcatDone.Error()) { - log.Debug(ctx, "got ErrConcatDone, exiting") + if err.Error() == fmt.Sprintf("%s: %s", ErrPipelineDone.Error(), ErrPipelineDone.Error()) { + log.Debug(ctx, "got ErrPipelineDone, exiting") return nil } log.Error(ctx, "gstreamer error", "error", err.Error()) diff --git a/pkg/media/concat2.go b/pkg/media/concat2.go deleted file mode 100644 index ffb3908dc..000000000 --- a/pkg/media/concat2.go +++ /dev/null @@ -1,258 +0,0 @@ -package media - -import ( - "context" - "errors" - "fmt" - - "github.com/go-gst/go-gst/gst" - "stream.place/streamplace/pkg/bus" - "stream.place/streamplace/pkg/log" -) - -var ErrConcatDone = errors.New("concat done") - -func ConcatBin(ctx context.Context, segCh <-chan *bus.Seg, doH264Parse bool) (*gst.Bin, error) { - ctx = log.WithLogValues(ctx, "func", "ConcatBin") - bin := gst.NewBin("concat-bin") - - streamsynchronizer, err := gst.NewElementWithProperties("streamsynchronizer", map[string]any{ - "name": "concat-streamsynchronizer", - }) - if err != nil { - return nil, fmt.Errorf("failed to create streamsynchronizer element: %w", err) - } - - err = bin.Add(streamsynchronizer) - if err != nil { - return nil, fmt.Errorf("failed to add streamsynchronizer to pipeline: %w", err) - } - - syncPadVideoSink := streamsynchronizer.GetRequestPad("sink_%u") - if syncPadVideoSink == nil { - return nil, fmt.Errorf("failed to get sync video sink pad") - } - - syncPadAudioSink := streamsynchronizer.GetRequestPad("sink_%u") - if syncPadAudioSink == nil { - return nil, fmt.Errorf("failed to get sync audio sink pad") - } - - syncPadVideoSrc := streamsynchronizer.GetStaticPad("src_0") - if syncPadVideoSrc == nil { - return nil, fmt.Errorf("failed to get sync video src pad") - } - - syncPadAudioSrc := streamsynchronizer.GetStaticPad("src_1") - if syncPadAudioSrc == nil { - return nil, fmt.Errorf("failed to get sync audio src pad") - } - - mq, err := gst.NewElementWithProperties("multiqueue", map[string]interface{}{ - "name": "concat-multiqueue", - }) - if err != nil { - return nil, fmt.Errorf("failed to create multiqueue element: %w", err) - } - err = bin.Add(mq) - if err != nil { - return nil, fmt.Errorf("failed to add multiqueue to bin: %w", err) - } - - // 10x default multiqueue size - err = mq.SetProperty("max-size-time", uint64(200000000000)) - if err != nil { - return nil, fmt.Errorf("failed to set max-size-time: %w", err) - } - err = mq.SetProperty("max-size-bytes", uint(1048576000)) - if err != nil { - return nil, fmt.Errorf("failed to set max-size-bytes: %w", err) - } - err = mq.SetProperty("max-size-buffers", uint(500)) - if err != nil { - return nil, fmt.Errorf("failed to set max-size-buffers: %w", err) - } - - mqVideoSink := mq.GetRequestPad("sink_%u") - if mqVideoSink == nil { - return nil, fmt.Errorf("video sink pad not found") - } - - mqAudioSink := mq.GetRequestPad("sink_%u") - if mqAudioSink == nil { - return nil, fmt.Errorf("audio sink pad not found") - } - - mqVideoSrc := mq.GetStaticPad("src_0") - if mqVideoSrc == nil { - return nil, fmt.Errorf("video source pad not found") - } - - mqAudioSrc := mq.GetStaticPad("src_1") - if mqAudioSrc == nil { - return nil, fmt.Errorf("audio source pad not found") - } - - linked := syncPadVideoSrc.Link(mqVideoSink) - if linked != gst.PadLinkOK { - return nil, fmt.Errorf("failed to link sync video src pad to multiqueue video sink pad: %v", linked) - } - - linked = syncPadAudioSrc.Link(mqAudioSink) - if linked != gst.PadLinkOK { - return nil, fmt.Errorf("failed to link sync audio src pad to multiqueue audio sink pad: %v", linked) - } - - videoGhost := gst.NewGhostPad("video_0", mqVideoSrc) - if videoGhost == nil { - return nil, fmt.Errorf("failed to create video ghost pad") - } - - audioGhost := gst.NewGhostPad("audio_0", mqAudioSrc) - if audioGhost == nil { - return nil, fmt.Errorf("failed to create audio ghost pad") - } - - ok := bin.AddPad(videoGhost.Pad) - if !ok { - return nil, fmt.Errorf("failed to add video ghost pad to bin") - } - - ok = bin.AddPad(audioGhost.Pad) - if !ok { - return nil, fmt.Errorf("failed to add audio ghost pad to bin") - } - - go func() { - for { - select { - case seg := <-segCh: - if seg == nil { - - ok := syncPadVideoSrc.PushEvent(gst.NewEOSEvent()) - if !ok { - log.Error(ctx, "failed to post EOS message", "error", ok) - } - ok = syncPadAudioSrc.PushEvent(gst.NewEOSEvent()) - if !ok { - log.Error(ctx, "failed to post EOS message", "error", ok) - } - log.Debug(ctx, "concat completed") - - return - } - err := addConcatDemuxer(ctx, bin, seg, syncPadVideoSink, syncPadAudioSink, doH264Parse) - if err != nil { - if ctx.Err() != nil { - // Session ended mid-segment; not a pipeline error. - return - } - log.Error(ctx, "failed to add concat demuxer", "error", err) - bin.Error(err.Error(), err) - return - } - case <-ctx.Done(): - return - } - } - }() - - return bin, nil -} - -func addConcatDemuxer(ctx context.Context, bin *gst.Bin, seg *bus.Seg, syncPadVideoSink *gst.Pad, syncPadAudioSink *gst.Pad, doH264Parse bool) error { - var cancel context.CancelFunc - ctx, cancel = context.WithCancel(ctx) - defer cancel() - ctx = log.WithLogValues(ctx, "func", "ConcatBin") - - log.Debug(ctx, "adding concat demuxer", "seg", seg.Filepath) - demuxBin, err := ConcatDemuxBin(ctx, seg, doH264Parse) - if err != nil { - return fmt.Errorf("failed to create demux bin: %w", err) - } - - err = bin.Add(demuxBin.Element) - if err != nil { - return fmt.Errorf("failed to add demux bin to bin: %w", err) - } - - demuxBinPadVideoSrc := demuxBin.GetStaticPad("video_0") - if demuxBinPadVideoSrc == nil { - return fmt.Errorf("failed to get demux bin video src pad") - } - - demuxBinPadAudioSrc := demuxBin.GetStaticPad("audio_0") - if demuxBinPadAudioSrc == nil { - return fmt.Errorf("failed to get demux bin audio src pad") - } - - linked := demuxBinPadVideoSrc.Link(syncPadVideoSink) - if linked != gst.PadLinkOK { - return fmt.Errorf("failed to link demux bin video src pad to sync video sink pad: %v", linked) - } - - linked = demuxBinPadAudioSrc.Link(syncPadAudioSink) - if linked != gst.PadLinkOK { - return fmt.Errorf("failed to link demux bin audio src pad to sync audio sink pad: %v", linked) - } - - // Buffered so the probe never blocks: if the context is cancelled - // below we stop receiving, but a late EOS must still be able to land - // without stranding the probe's calling thread. - eosCh := make(chan struct{}, 2) - eos := func(pad *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { - if pad.GetDirection() != gst.PadDirectionSource { - return gst.PadProbeOK - } - if info.GetEvent().Type() != gst.EventTypeEOS { - return gst.PadProbeOK - } - log.Debug(ctx, "demux EOS", "name", pad.GetName(), "direction", pad.GetDirection()) - downstreamPad := pad.GetPeer() - unlinked := pad.Unlink(downstreamPad) - if !unlinked { - log.Error(ctx, "failed to unlink pad", "name", pad.GetName(), "direction", pad.GetDirection(), "error", unlinked) - } - eosCh <- struct{}{} - return gst.PadProbeRemove - } - demuxBinPadVideoSrc.AddProbe(gst.PadProbeTypeEventBoth, eos) - demuxBinPadAudioSrc.AddProbe(gst.PadProbeTypeEventBoth, eos) - - if err := bin.SetState(gst.StatePlaying); err != nil { - return fmt.Errorf("failed to set state: %w", err) - } - - // Wait for EOS on both demux src pads. A segment that never emits EOS - // on both pads (audio-only, malformed, or a session torn down before - // the stream completes) would otherwise wedge this goroutine — and the - // demux bin's GStreamer state — forever. Bail out on cancellation. - for range 2 { - select { - case <-eosCh: - case <-ctx.Done(): - if rmErr := bin.Remove(demuxBin.Element); rmErr != nil { - log.Error(ctx, "failed to remove demux bin after cancel", "error", rmErr) - } - if stErr := demuxBin.SetState(gst.StateNull); stErr != nil { - log.Error(ctx, "failed to set demux bin to null after cancel", "error", stErr) - } - return ctx.Err() - } - } - - err = bin.Remove(demuxBin.Element) - if err != nil { - return fmt.Errorf("failed to remove demux bin from bin: %w", err) - } - - err = demuxBin.SetState(gst.StateNull) - if err != nil { - return fmt.Errorf("failed to set demux bin to null state: %w", err) - } - - log.Debug(ctx, "removed concat demuxer", "seg", seg.Filepath) - - return nil -} diff --git a/pkg/media/concat2_test.go b/pkg/media/concat2_test.go deleted file mode 100644 index 9adb999e5..000000000 --- a/pkg/media/concat2_test.go +++ /dev/null @@ -1,294 +0,0 @@ -package media - -import ( - "bytes" - "context" - "fmt" - "io" - "os" - "testing" - "time" - - "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" - "github.com/google/uuid" - "github.com/stretchr/testify/require" - "golang.org/x/sync/errgroup" - "stream.place/streamplace/pkg/bus" - "stream.place/streamplace/pkg/log" -) - -// TestAddConcatDemuxerUnblocksOnCancel is a regression test for a -// goroutine leak: addConcatDemuxer waits for EOS on both the video and -// audio demux src pads before returning. A segment that only ever emits -// EOS on one pad (audio-only, malformed, or a session torn down before -// the stream completes) used to wedge the goroutine — and the demux -// bin's GStreamer state — forever, with no way for the parent context -// to interrupt it. Over a long-lived server these accumulated until a -// restart. -// -// We feed the audio-only fixture so exactly one EOS ever fires, then -// cancel the context. The fix makes addConcatDemuxer observe -// cancellation and return; before the fix this blocks forever and the -// test trips its deadline. -func TestAddConcatDemuxerUnblocksOnCancel(t *testing.T) { - ctx := log.WithLogValues(context.Background(), "test", "TestAddConcatDemuxerUnblocksOnCancel") - ctx, cancel := context.WithCancel(ctx) - defer cancel() - - pipeline, err := gst.NewPipeline("TestAddConcatDemuxerUnblocksOnCancel") - require.NoError(t, err) - - // Minimal stand-in for the concat bin's wiring: a streamsynchronizer - // whose two request sink pads are the link points addConcatDemuxer - // expects, with fakesinks downstream so data actually flows. - bin := gst.NewBin("concat-bin") - streamsync, err := gst.NewElementWithProperties("streamsynchronizer", map[string]any{"name": "ss"}) - require.NoError(t, err) - require.NoError(t, bin.Add(streamsync)) - - videoSink := streamsync.GetRequestPad("sink_%u") - require.NotNil(t, videoSink) - audioSink := streamsync.GetRequestPad("sink_%u") - require.NotNil(t, audioSink) - - for i, srcName := range []string{"src_0", "src_1"} { - fakesink, err := gst.NewElementWithProperties("fakesink", map[string]any{ - "name": fmt.Sprintf("fakesink_%d", i), - "sync": false, - }) - require.NoError(t, err) - require.NoError(t, bin.Add(fakesink)) - require.Equal(t, gst.PadLinkOK, - streamsync.GetStaticPad(srcName).Link(fakesink.GetStaticPad("sink"))) - } - - require.NoError(t, pipeline.Add(bin.Element)) - - go func() { _ = HandleBusMessages(ctx, pipeline) }() - require.NoError(t, pipeline.SetState(gst.StatePlaying)) - defer func() { _ = pipeline.BlockSetState(gst.StateNull) }() - - data, err := os.ReadFile(getFixture("duration-mismatch-audio.mp4")) - require.NoError(t, err) - seg := &bus.Seg{Data: data, Filepath: "duration-mismatch-audio.mp4"} - - done := make(chan error, 1) - go func() { - done <- addConcatDemuxer(ctx, bin, seg, videoSink, audioSink, true) - }() - - // Let the single audio EOS arrive and be consumed, leaving - // addConcatDemuxer blocked on the second (video) EOS that never comes. - time.Sleep(2 * time.Second) - cancel() - - select { - case err := <-done: - require.ErrorIs(t, err, context.Canceled) - case <-time.After(15 * time.Second): - t.Fatal("addConcatDemuxer did not return after context cancellation — eosCh leak still present") - } -} - -func TestConcatBin(t *testing.T) { - withNoGSTLeaks(t, func() { - - g, _ := errgroup.WithContext(context.Background()) - for range streamplaceTestCount { - g.Go(func() error { - return innerTestConcatBin(t) - }) - } - err := g.Wait() - require.NoError(t, err) - }) -} - -// This function remains in scope for the duration of a single users' playback -func innerTestConcatBin(t *testing.T) error { - ctx := log.WithDebugValue(context.Background(), map[string]map[string]int{"func": {"ConcatStream": 9, "ConcatBin": 9, "SegDemuxBin": 9}}) - tag := os.Getenv("TEST_TAG") - uuid, _ := uuid.NewV7() - uuidStr := uuid.String() - if tag != "" { - ctx = log.WithLogValues(ctx, "tag", tag) - uuidStr = fmt.Sprintf("%s-%s", tag, uuidStr) - } - ctx = log.WithLogValues(ctx, "func", "ConcatBin", "uuid", uuidStr) - - pipeline, err := gst.NewPipeline("TestConcatBin") - if err != nil { - return fmt.Errorf("failed to create pipeline: %w", err) - } - - ctx, cancel := context.WithCancel(ctx) - - errCh := make(chan error) - go func() { - err := HandleBusMessages(ctx, pipeline) - cancel() - errCh <- err - close(errCh) - }() - - defer func() { - cancel() - err := <-errCh - require.NoError(t, err, fmt.Sprintf("uuid: %s", uuidStr)) - err = pipeline.BlockSetState(gst.StateNull) - require.NoError(t, err, fmt.Sprintf("uuid: %s", uuidStr)) - }() - - filename := getFixture("sample-segment.mp4") - inputFile, err := os.Open(filename) - if err != nil { - return fmt.Errorf("failed to open fixture file: %w", err) - } - defer inputFile.Close() - - bs, err := io.ReadAll(inputFile) - if err != nil { - return fmt.Errorf("failed to read fixture file: %w", err) - } - - testSegs := []*bus.Seg{} - for range 5 { - testSegs = append(testSegs, &bus.Seg{ - Data: bs, - Filepath: filename, - }) - } - - segCh := make(chan *bus.Seg) - go func() { - for _, seg := range testSegs { - segCh <- seg - } - close(segCh) - }() - - concatBin, err := ConcatBin(ctx, segCh, true) - if err != nil { - return fmt.Errorf("failed to create concat bin: %w", err) - } - - err = pipeline.Add(concatBin.Element) - if err != nil { - return fmt.Errorf("failed to add concat bin to pipeline: %w", err) - } - - videoPad := concatBin.GetStaticPad("video_0") - if videoPad == nil { - return fmt.Errorf("video pad not found") - } - - audioPad := concatBin.GetStaticPad("audio_0") - if audioPad == nil { - return fmt.Errorf("audio pad not found") - } - - videoAppSink, err := gst.NewElementWithProperties("appsink", map[string]interface{}{ - "name": "videoappsink", - "sync": false, - }) - if err != nil { - return fmt.Errorf("failed to create video appsink: %w", err) - } - - err = pipeline.Add(videoAppSink) - if err != nil { - return fmt.Errorf("failed to add video appsink to pipeline: %w", err) - } - - videoAppSinkPadSink := videoAppSink.GetStaticPad("sink") - if videoAppSinkPadSink == nil { - return fmt.Errorf("video appsink pad not found") - } - - audioAppSink, err := gst.NewElementWithProperties("appsink", map[string]interface{}{ - "name": "audioappsink", - "sync": false, - }) - if err != nil { - return fmt.Errorf("failed to create audio appsink: %w", err) - } - - err = pipeline.Add(audioAppSink) - if err != nil { - return fmt.Errorf("failed to add audio appsink to pipeline: %w", err) - } - - audioAppSinkPadSink := audioAppSink.GetStaticPad("sink") - if audioAppSinkPadSink == nil { - return fmt.Errorf("audio appsink pad not found") - } - - ok := videoPad.Link(videoAppSinkPadSink) - if ok != gst.PadLinkOK { - return fmt.Errorf("failed to link video pad: %v", ok) - } - - ok = audioPad.Link(audioAppSinkPadSink) - if ok != gst.PadLinkOK { - return fmt.Errorf("failed to link audio pad: %v", ok) - } - - videoBuf := bytes.Buffer{} - audioBuf := bytes.Buffer{} - - videoappsink := app.SinkFromElement(videoAppSink) - videoappsink.SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: WriterNewSample(ctx, &videoBuf), - }) - - audioappsink := app.SinkFromElement(audioAppSink) - audioappsink.SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: WriterNewSample(ctx, &audioBuf), - }) - - // Start the pipeline - err = pipeline.SetState(gst.StatePlaying) - if err != nil { - return fmt.Errorf("failed to set pipeline to playing state: %w", err) - } - - // Start a goroutine to print buffer sizes - go func() { - for { - select { - case <-ctx.Done(): - return - case <-time.After(1 * time.Second): - log.Debug(ctx, "buffer sizes", - "videoBuf", videoBuf.Len(), - "audioBuf", audioBuf.Len()) - } - } - }() - - <-ctx.Done() - - time.Sleep(5 * time.Second) - - padIdleCh := make(chan struct{}) - - padIdle := func(pad *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { - log.Debug(ctx, "pad-idle", "name", pad.GetName(), "direction", pad.GetDirection()) - go func() { - padIdleCh <- struct{}{} - }() - return gst.PadProbeRemove - } - - videoAppSinkPadSink.AddProbe(gst.PadProbeTypeIdle, padIdle) - audioAppSinkPadSink.AddProbe(gst.PadProbeTypeIdle, padIdle) - - <-padIdleCh - <-padIdleCh - - require.Equal(t, 4936455, videoBuf.Len(), fmt.Sprintf("uuid: %s", uuidStr)) - require.Equal(t, 32200, audioBuf.Len(), fmt.Sprintf("uuid: %s", uuidStr)) - - return <-errCh -} diff --git a/pkg/media/thumbnail.go b/pkg/media/thumbnail.go index 8a4570cee..1dd5a9986 100644 --- a/pkg/media/thumbnail.go +++ b/pkg/media/thumbnail.go @@ -125,7 +125,7 @@ func thumbnailFromMP4(ctx context.Context, flat []byte, w io.Writer, format stri select { case <-thumbCh: // Got the frame; wind the pipeline down so HandleBusMessages returns. - pipeline.Error(ErrConcatDone.Error(), ErrConcatDone) + pipeline.Error(ErrPipelineDone.Error(), ErrPipelineDone) <-errCh return nil case err := <-errCh: @@ -272,7 +272,7 @@ func Thumbnail(ctx context.Context, r io.Reader, w io.Writer, format string) err <-thumbCh // signals the pipeline to clean up cleanly - pipeline.Error(ErrConcatDone.Error(), ErrConcatDone) + pipeline.Error(ErrPipelineDone.Error(), ErrPipelineDone) busErr := <-errCh