diff --git a/pkg/bus/segchanman.go b/pkg/bus/segchanman.go index 06c28a72e..570a6d4de 100644 --- a/pkg/bus/segchanman.go +++ b/pkg/bus/segchanman.go @@ -20,9 +20,20 @@ type Seg struct { Published bool } +// PacketizedSample is one WebRTC-writable sample — a video access unit or an +// Opus packet — with its real duration on the source timeline. Carrying the +// per-sample duration (rather than dividing the segment evenly) keeps +// non-uniform frame spacing intact: an encoder shedding frames under +// bandwidth pressure sends bursts and gaps, and respacing those uniformly +// rubber-bands the video against the audio. +type PacketizedSample struct { + Data []byte + Duration time.Duration +} + type PacketizedSegment struct { - Video [][]byte - Audio [][]byte + Video []PacketizedSample + Audio []PacketizedSample Duration time.Duration } diff --git a/pkg/media/packetize.go b/pkg/media/packetize.go index 99ecfe422..843cfc868 100644 --- a/pkg/media/packetize.go +++ b/pkg/media/packetize.go @@ -138,8 +138,8 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti return nil, fmt.Errorf("failed to get audio appsink element") } - videoOutput := [][]byte{} - audioOutput := [][]byte{} + videoOutput := []rawSample{} + audioOutput := []rawSample{} // eosCh := make(chan struct{}) // Slice-less video data (parameter sets / SEI metadata with no picture) @@ -178,7 +178,7 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti pendingVideo = nil } - videoOutput = append(videoOutput, samples) + videoOutput = append(videoOutput, newRawSample(samples, buffer)) return gst.FlowOK }, @@ -206,7 +206,7 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti samples := buffer.Bytes() // log.Warn(ctx, "audioappsink NewSampleFunc", "sample", len(samples)) - audioOutput = append(audioOutput, samples) + audioOutput = append(audioOutput, newRawSample(samples, buffer)) clockTime := buffer.Duration() dur := clockTime.AsDuration() @@ -255,8 +255,58 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti } return &bus.PacketizedSegment{ - Video: videoOutput, - Audio: audioOutput, + Video: finalizeSampleDurations(videoOutput, segDur), + Audio: finalizeSampleDurations(audioOutput, segDur), Duration: segDur, }, nil } + +// rawSample is a demuxed sample plus the source timing needed to compute its +// WebRTC duration: its decode timestamp (DTS when the container carries one — +// decode order is monotonic even with B-frames — else PTS) and the buffer's +// own duration as a fallback. +type rawSample struct { + data []byte + ts time.Duration + hasTS bool + bufDur time.Duration + hasDur bool +} + +func newRawSample(data []byte, buffer *gst.Buffer) rawSample { + rs := rawSample{data: data} + if ts := buffer.DecodingTimestamp().AsDuration(); ts != nil { + rs.ts, rs.hasTS = *ts, true + } else if pts := buffer.PresentationTimestamp().AsDuration(); pts != nil { + rs.ts, rs.hasTS = *pts, true + } + if dur := buffer.Duration().AsDuration(); dur != nil { + rs.bufDur, rs.hasDur = *dur, true + } + return rs +} + +// finalizeSampleDurations converts raw buffer timing into the per-sample +// durations the WebRTC sender stamps and paces by. Each sample lasts until +// the next sample's timestamp — the source timeline, so non-uniform spacing +// (an encoder shedding frames under bandwidth pressure) is preserved instead +// of being respaced evenly. The last sample has no successor: it gets its own +// buffer duration, stretched to fill out the segment (total) when the track +// would otherwise end early — a sparse track's final frame must hold until +// the segment ends or its timeline falls behind the other track's. +func finalizeSampleDurations(raw []rawSample, total time.Duration) []bus.PacketizedSample { + out := make([]bus.PacketizedSample, len(raw)) + var span time.Duration + for i, rs := range raw { + dur := rs.bufDur + if i+1 < len(raw) && rs.hasTS && raw[i+1].hasTS && raw[i+1].ts > rs.ts { + dur = raw[i+1].ts - rs.ts + } + out[i] = bus.PacketizedSample{Data: rs.data, Duration: dur} + span += dur + } + if len(out) > 0 && total > span { + out[len(out)-1].Duration += total - span + } + return out +} diff --git a/pkg/media/packetize_test.go b/pkg/media/packetize_test.go index a96fe4ee7..576fefafe 100644 --- a/pkg/media/packetize_test.go +++ b/pkg/media/packetize_test.go @@ -266,8 +266,15 @@ func TestPacketizeTrailingCaptionSEI(t *testing.T) { require.Equal(t, frames, len(packet.Video)) totalSEIs := 0 for i, v := range packet.Video { - require.True(t, hasVideoSlice(v), "video sample %d has no picture", i) - totalSEIs += countCaptionSEIs(v) + require.True(t, hasVideoSlice(v.Data), "video sample %d has no picture", i) + totalSEIs += countCaptionSEIs(v.Data) + // Real per-sample timing: ~33ms at 30fps, no sample burned by a + // zero-length caption "frame". The last sample is exempt — it + // stretches to cover the (audio-derived) segment end. + if i < len(packet.Video)-1 { + require.InDelta(t, 33*time.Millisecond, v.Duration, float64(10*time.Millisecond), + "video sample %d duration", i) + } } // ...and the captions must survive, riding with real frames. require.Equal(t, seiCount, totalSEIs) @@ -275,6 +282,64 @@ func TestPacketizeTrailingCaptionSEI(t *testing.T) { }) } +// TestFinalizeSampleDurations covers the duration synthesis rules: durations +// come from successive timestamps (preserving gaps from dropped frames), the +// buffer's own duration is the fallback, and the final sample stretches to +// the end of the segment when the track would otherwise end early. +func TestFinalizeSampleDurations(t *testing.T) { + ms := func(n int) time.Duration { return time.Duration(n) * time.Millisecond } + + t.Run("UniformTimeline", func(t *testing.T) { + raw := []rawSample{ + {ts: ms(0), hasTS: true, bufDur: ms(33), hasDur: true}, + {ts: ms(33), hasTS: true, bufDur: ms(33), hasDur: true}, + {ts: ms(66), hasTS: true, bufDur: ms(34), hasDur: true}, + } + out := finalizeSampleDurations(raw, ms(100)) + require.Equal(t, []time.Duration{ms(33), ms(33), ms(34)}, + []time.Duration{out[0].Duration, out[1].Duration, out[2].Duration}) + }) + + t.Run("GapFromDroppedFrames", func(t *testing.T) { + // An encoder under bandwidth pressure sent two frames, dropped ~1s, + // then sent another: the gap belongs to the sample before it. + raw := []rawSample{ + {ts: ms(0), hasTS: true, bufDur: ms(33), hasDur: true}, + {ts: ms(33), hasTS: true, bufDur: ms(33), hasDur: true}, + {ts: ms(1033), hasTS: true, bufDur: ms(33), hasDur: true}, + } + out := finalizeSampleDurations(raw, ms(1066)) + require.Equal(t, ms(33), out[0].Duration) + require.Equal(t, ms(1000), out[1].Duration) + require.Equal(t, ms(33), out[2].Duration) + }) + + t.Run("LastSampleStretchesToSegmentEnd", func(t *testing.T) { + // A single keyframe in a 4s segment must hold the full 4s, or the + // video timeline falls behind the audio's segment after segment. + raw := []rawSample{ + {ts: ms(0), hasTS: true, bufDur: ms(33), hasDur: true}, + } + out := finalizeSampleDurations(raw, ms(4000)) + require.Equal(t, ms(4000), out[0].Duration) + }) + + t.Run("NoStretchWhenTrackFillsSegment", func(t *testing.T) { + // Video already spans past the (audio-derived) total: leave it alone. + raw := []rawSample{ + {ts: ms(0), hasTS: true, bufDur: ms(33), hasDur: true}, + {ts: ms(33), hasTS: true, bufDur: ms(34), hasDur: true}, + } + out := finalizeSampleDurations(raw, ms(50)) + require.Equal(t, ms(33), out[0].Duration) + require.Equal(t, ms(34), out[1].Duration) + }) + + t.Run("Empty", func(t *testing.T) { + require.Empty(t, finalizeSampleDurations(nil, ms(1000))) + }) +} + func TestPacketizeInvalid(t *testing.T) { // cur := goleak.IgnoreCurrent() // defer goleak.VerifyNone(t, cur) diff --git a/pkg/media/webrtc_playback2.go b/pkg/media/webrtc_playback2.go index b17a217ea..e287eb6e2 100644 --- a/pkg/media/webrtc_playback2.go +++ b/pkg/media/webrtc_playback2.go @@ -161,63 +161,30 @@ func (mm *MediaManager) WebRTCPlayback2(ctx context.Context, user string, rendit latency -= packet.Duration scalar = getPlaybackRate(latency) log.Debug(ctx, "playback latency", "latency", latency, "scalar", scalar) - var videoDur time.Duration - var audioDur time.Duration - if len(packet.Video) > 0 { - videoDur = packet.Duration / time.Duration(len(packet.Video)) - } - if len(packet.Audio) > 0 { - audioDur = packet.Duration / time.Duration(len(packet.Audio)) - } g, _ := errgroup.WithContext(ctx) + wroteAny := false - if !audioOnly && videoDur > 0 { + if !audioOnly && len(packet.Video) > 0 { + wroteAny = true g.Go(func() error { - ticker := time.NewTicker(time.Duration(float64(videoDur) * (1 / scalar))) - defer ticker.Stop() - for _, video := range packet.Video { - err := videoTrack.WriteSample(media.Sample{Data: video, Duration: videoDur}) - if err != nil { - return fmt.Errorf("failed to write video sample: %w", err) - } - - select { - case <-ctx.Done(): - return nil - case <-ticker.C: - continue - } - } - return nil + return writeSamples(ctx, videoTrack, packet.Video, scalar) }) } else if !audioOnly { log.Warn(ctx, "no video samples to write") } - if audioDur > 0 { + if len(packet.Audio) > 0 { + wroteAny = true g.Go(func() error { - ticker := time.NewTicker(time.Duration(float64(audioDur) * (1 / scalar))) - defer ticker.Stop() - for _, audio := range packet.Audio { - err := audioTrack.WriteSample(media.Sample{Data: audio, Duration: audioDur}) - if err != nil { - return fmt.Errorf("failed to write audio sample: %w", err) - } - select { - case <-ctx.Done(): - return nil - case <-ticker.C: - continue - } - } - return nil + return writeSamples(ctx, audioTrack, packet.Audio, scalar) }) - + } else { + log.Warn(ctx, "no audio samples to write") + } + if wroteAny { if err := g.Wait(); err != nil { log.Error(ctx, "failed to write samples", "error", err) cancel() } - } else { - log.Warn(ctx, "no audio samples to write") } } } @@ -280,6 +247,32 @@ func (mm *MediaManager) WebRTCPlayback2(ctx context.Context, user string, rendit } } +// writeSamples writes one track's samples paced by their real durations — +// the source timeline, so non-uniform frame spacing (bursts and gaps from an +// encoder shedding frames) reaches the viewer intact. The stamped duration +// stays real even while catch-up (scalar > 1) speeds the pacing: the +// receiver's clock is authoritative for playout, sending faster just refills +// its buffer. +func writeSamples(ctx context.Context, track *webrtc.TrackLocalStaticSample, samples []bus.PacketizedSample, scalar float64) error { + for _, s := range samples { + if err := track.WriteSample(media.Sample{Data: s.Data, Duration: s.Duration}); err != nil { + return fmt.Errorf("failed to write sample: %w", err) + } + wait := time.Duration(float64(s.Duration) / scalar) + if wait <= 0 { + continue + } + timer := time.NewTimer(wait) + select { + case <-ctx.Done(): + timer.Stop() + return nil + case <-timer.C: + } + } + return nil +} + // getPlaybackRate returns a playback rate that eases from 1.0 to 1.5 between 7 and 60 seconds func getPlaybackRate(dur time.Duration) float64 { switch {