diff --git a/pkg/media/concat_demux.go b/pkg/media/concat_demux.go index 44b0170fc..42cb5b527 100644 --- a/pkg/media/concat_demux.go +++ b/pkg/media/concat_demux.go @@ -5,6 +5,7 @@ import ( "context" "fmt" "strings" + "sync/atomic" "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" @@ -176,9 +177,16 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B 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 - padAdded = func(self *gst.Element, pad *gst.Pad) { + // The signal handlers below are swapped for no-ops once they've done + // their job (or the ctx ends, whichever comes first) so their closures + // release the mq pads — and, via the pads' sticky events, tag lists and + // samples — that they capture; without the swap those leak. The swap + // happens on the ctx goroutine while the demux's dispatch thread reads + // the same slot, so both sides go through an atomic (a plain assignment + // here is a data race; a mutex held across dispatch risks blocking the + // swap behind a stuck handler, pinning the closures). + var padAdded atomic.Value + padAdded.Store(func(self *gst.Element, pad *gst.Pad) { var downstreamPad *gst.Pad if strings.HasPrefix(pad.GetName(), "video_") { downstreamPad = mqVideoSink @@ -199,24 +207,11 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B } needed-- if needed == 0 { - padAdded = doNothing + padAdded.Store(doNothing) } - } + }) outerPadAdded := func(self *gst.Element, pad *gst.Pad) { - padAdded(self, pad) - } - - // Necessary to avoid leaking `mqVideoSink` and `mqAudioSink` from the - // pad-added function in the case where we hit invalid data and - // pad-added never fires. - go func() { - <-ctx.Done() - padAdded = doNothing - }() - - _, err = demux.Connect("pad-added", outerPadAdded) - if err != nil { - return nil, fmt.Errorf("failed to connect demux pad-added signal: %w", err) + padAdded.Load().(func(self *gst.Element, pad *gst.Pad))(self, pad) } // Both branches are wired up before the demux runs, but the segment may @@ -228,10 +223,9 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B // 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) { + // linked flags need no synchronization. + var noMorePads atomic.Value + noMorePads.Store(func(self *gst.Element) { if !videoLinked { log.Warn(ctx, "segment has no video track; completing video branch with EOS") mqVideoSink.SendEvent(gst.NewEOSEvent()) @@ -240,15 +234,24 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B log.Warn(ctx, "segment has no audio track; completing audio branch with EOS") mqAudioSink.SendEvent(gst.NewEOSEvent()) } - noMorePads = doNothingElement - } + noMorePads.Store(doNothingElement) + }) outerNoMorePads := func(self *gst.Element) { - noMorePads(self) + noMorePads.Load().(func(self *gst.Element))(self) } + + // Necessary to avoid leaking `mqVideoSink` and `mqAudioSink` from the + // handlers in the case where we hit invalid data and they never fire. go func() { <-ctx.Done() - noMorePads = doNothingElement + padAdded.Store(doNothing) + noMorePads.Store(doNothingElement) }() + + _, err = demux.Connect("pad-added", outerPadAdded) + if err != nil { + return nil, fmt.Errorf("failed to connect demux pad-added signal: %w", err) + } _, err = demux.Connect("no-more-pads", outerNoMorePads) if err != nil { return nil, fmt.Errorf("failed to connect demux no-more-pads signal: %w", err) diff --git a/pkg/media/packetize.go b/pkg/media/packetize.go index 843cfc868..187cc8ab9 100644 --- a/pkg/media/packetize.go +++ b/pkg/media/packetize.go @@ -254,10 +254,20 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti return nil, fmt.Errorf("packetize pipeline error filename=%s, error=%w", seg.Filepath, err) } + video := finalizeSampleDurations(videoOutput, segDur) + // segDur is audio-derived; a video-only segment would report Duration 0, + // throwing off the sender's latency bookkeeping (the segment occupies the + // sender for its full video span). Fall back to the video timeline. + duration := segDur + if duration == 0 { + for _, s := range video { + duration += s.Duration + } + } return &bus.PacketizedSegment{ - Video: finalizeSampleDurations(videoOutput, segDur), + Video: video, Audio: finalizeSampleDurations(audioOutput, segDur), - Duration: segDur, + Duration: duration, }, nil } diff --git a/pkg/media/packetize_test.go b/pkg/media/packetize_test.go index 39d446be4..033c6862a 100644 --- a/pkg/media/packetize_test.go +++ b/pkg/media/packetize_test.go @@ -359,6 +359,10 @@ func TestPacketizeSingleTrackSegment(t *testing.T) { require.NoError(t, err) require.Equal(t, 30, len(packet.Video)) require.Empty(t, packet.Audio) + // Duration falls back to the video span when there's no audio to + // derive it from — the sender's latency accounting needs the time + // this segment actually occupies. 30 frames at 30fps ≈ 1s. + require.InDelta(t, time.Second, packet.Duration, float64(100*time.Millisecond)) }) }) t.Run("AudioOnly", func(t *testing.T) {