diff --git a/pkg/media/concat_demux.go b/pkg/media/concat_demux.go index c556d7e71..44b0170fc 100644 --- a/pkg/media/concat_demux.go +++ b/pkg/media/concat_demux.go @@ -16,6 +16,9 @@ import ( // silly technique to avoid leaking pads func doNothing(self *gst.Element, pad *gst.Pad) {} +// doNothing for signals without a pad argument +func doNothingElement(self *gst.Element) {} + // Function for demuxing a single segment. Needs to be handled very carefully. // In particular: users of this MUST cancel the passed context when they're // done with the bin. @@ -170,6 +173,8 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B } needed := 2 + videoLinked := false + audioLinked := false var padAdded func(self *gst.Element, pad *gst.Pad) // the defer funcs are needed to avoid leaking pads for some reason @@ -177,8 +182,10 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B var downstreamPad *gst.Pad if strings.HasPrefix(pad.GetName(), "video_") { downstreamPad = mqVideoSink + videoLinked = true } else if strings.HasPrefix(pad.GetName(), "audio_") { downstreamPad = mqAudioSink + audioLinked = true } else { log.Error(ctx, "unknown pad", "name", pad.GetName(), "direction", pad.GetDirection()) // cancel() @@ -212,6 +219,41 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B return nil, fmt.Errorf("failed to connect demux pad-added signal: %w", err) } + // Both branches are wired up before the demux runs, but the segment may + // not have both tracks: a stream whose encoder is shedding everything but + // keyframes under bandwidth pressure can produce a segment with video and + // no audio (or, in principle, the reverse). The demux only EOSes pads it + // created, so a branch whose pad never appears would never see data OR + // EOS — its sink never finishes, the pipeline never completes, and the + // caller hangs until its timeout. When the demux declares it's done making + // pads, flush any branch left orphaned with an EOS so the pipeline can + // drain. Fires on the demux's streaming thread, same as pad-added, so the + // linked flags need no locking. Uses the same swap-to-no-op indirection as + // padAdded above so the closure's pad references don't leak. + var noMorePads func(self *gst.Element) + noMorePads = func(self *gst.Element) { + if !videoLinked { + log.Warn(ctx, "segment has no video track; completing video branch with EOS") + mqVideoSink.SendEvent(gst.NewEOSEvent()) + } + if !audioLinked { + log.Warn(ctx, "segment has no audio track; completing audio branch with EOS") + mqAudioSink.SendEvent(gst.NewEOSEvent()) + } + noMorePads = doNothingElement + } + outerNoMorePads := func(self *gst.Element) { + noMorePads(self) + } + go func() { + <-ctx.Done() + noMorePads = doNothingElement + }() + _, err = demux.Connect("no-more-pads", outerNoMorePads) + if err != nil { + return nil, fmt.Errorf("failed to connect demux no-more-pads signal: %w", err) + } + ok := bin.AddPad(videoGhost.Pad) if !ok { return nil, fmt.Errorf("failed to add video ghost pad to bin") diff --git a/pkg/media/packetize_test.go b/pkg/media/packetize_test.go index 576fefafe..39d446be4 100644 --- a/pkg/media/packetize_test.go +++ b/pkg/media/packetize_test.go @@ -340,6 +340,41 @@ func TestFinalizeSampleDurations(t *testing.T) { }) } +// TestPacketizeSingleTrackSegment: a segment can arrive with only one track — +// notably video-with-no-audio, seen in the wild when an encoder under +// bandwidth pressure sheds everything but keyframes and a whole GoP's worth +// of audio goes missing. ConcatDemuxBin pre-wires both branches, and the +// demux only EOSes pads it actually created, so before the no-more-pads +// backstop the trackless branch never completed: Packetize hung until its +// timeout and the segment vanished from WebRTC playback entirely. It must +// complete promptly with the missing track empty instead. +func TestPacketizeSingleTrackSegment(t *testing.T) { + t.Run("VideoOnly", func(t *testing.T) { + withNoGSTLeaks(t, func() { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + flat := runSynthPipeline(t, ctx, + "videotestsrc num-buffers=30 ! video/x-raw,width=320,height=240,framerate=30/1 ! x264enc tune=zerolatency speed-preset=ultrafast ! h264parse ! mp4mux fragment-duration=500 ! appsink name=sink") + packet, err := Packetize(context.Background(), &config.CLI{}, &bus.Seg{Data: flat}) + require.NoError(t, err) + require.Equal(t, 30, len(packet.Video)) + require.Empty(t, packet.Audio) + }) + }) + t.Run("AudioOnly", func(t *testing.T) { + withNoGSTLeaks(t, func() { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + flat := runSynthPipeline(t, ctx, + "audiotestsrc num-buffers=48 samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! opusenc ! mp4mux fragment-duration=500 ! appsink name=sink") + packet, err := Packetize(context.Background(), &config.CLI{}, &bus.Seg{Data: flat}) + require.NoError(t, err) + require.Empty(t, packet.Video) + require.NotEmpty(t, packet.Audio) + }) + }) +} + func TestPacketizeInvalid(t *testing.T) { // cur := goleak.IgnoreCurrent() // defer goleak.VerifyNone(t, cur)