From 1da503f397b80a1664b6043d05a70b317b3245ba Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 30 Jul 2026 20:10:56 -0700 Subject: [PATCH] webrtc: pace and stamp playback with real per-sample durations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PacketizedSegment carried only frame bytes and a segment total, so the sender divided the total evenly across the frames: every sample got the same synthesized duration and the same tick. That respaces non-uniform video — an encoder shedding frames under bandwidth pressure sends bursts and gaps — into a smear: frames display at the wrong times, rubber-banding against the audio. Carry each demuxed sample's real duration instead, derived from successive decode timestamps (so gaps from dropped frames survive, and B-frame decode order stays monotonic), with the buffer's own duration as fallback. The final sample stretches to the segment's end when the track would otherwise finish early, so a single-keyframe segment holds its frame for the full segment rather than letting the video timeline fall behind the audio's. The sender now also waits for both tracks on every segment (previously a segment with no audio left the video writer running unwaited, overlapping the next segment's writes), and catch-up speed only affects pacing — stamped durations stay real, since the receiver's clock is authoritative for playout. Co-Authored-By: Claude Fable 5 --- pkg/bus/segchanman.go | 15 ++++++- pkg/media/packetize.go | 62 ++++++++++++++++++++++++--- pkg/media/packetize_test.go | 69 ++++++++++++++++++++++++++++- pkg/media/webrtc_playback2.go | 81 ++++++++++++++++------------------- 4 files changed, 173 insertions(+), 54 deletions(-) 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 { -- 2.51.2