From 0ecb473fbed6a88b955691f71cb5180bd7bb6450 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 30 Jul 2026 20:04:58 -0700 Subject: [PATCH 1/5] webrtc: don't emit slice-less caption-SEI AUs as video frames MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Streams with embedded closed captions (e.g. OBS's caption API feeding a live transcription) carry a trailing CEA-708 caption SEI after some frames' slices. When Packetize's h264parse re-parses the byte stream it splits that SEI into its own slice-less, timestamp-less AU, which then got counted and sent to WebRTC viewers as a standalone video "frame": strict decoders (iOS VideoToolbox) error on a picture-less access unit, and since playback ignores PLI the picture stayed broken until the next keyframe — with long-GoP sources, seconds away. The bogus frames also inflated the frame count that the sender's synthesized per-frame timing divides by. Fold slice-less buffers into the next real frame instead — a leading SEI is where captions normally live, so caption-aware receivers still get them. A remainder at segment end has no frame to ride with and is dropped. Co-Authored-By: Claude Fable 5 --- pkg/media/packetize.go | 55 ++++++++-- pkg/media/packetize_test.go | 204 ++++++++++++++++++++++++++++++++++++ 2 files changed, 251 insertions(+), 8 deletions(-) diff --git a/pkg/media/packetize.go b/pkg/media/packetize.go index dd5906f21..99ecfe422 100644 --- a/pkg/media/packetize.go +++ b/pkg/media/packetize.go @@ -17,6 +17,32 @@ import ( "stream.place/streamplace/pkg/muxl" ) +// hasVideoSlice reports whether byte-stream H264 data contains a VCL NAL +// (coded slice, nal_unit_type 1–5) — i.e. an actual picture, as opposed to +// bare parameter sets or SEI metadata. +func hasVideoSlice(data []byte) bool { + for i := 0; i+3 < len(data); i++ { + if data[i] != 0 || data[i+1] != 0 { + continue + } + var nalIdx int + if data[i+2] == 1 { + nalIdx = i + 3 + } else if data[i+2] == 0 && i+4 < len(data) && data[i+3] == 1 { + nalIdx = i + 4 + } else { + continue + } + if nalIdx < len(data) { + if t := data[nalIdx] & 0x1f; t >= 1 && t <= 5 { + return true + } + } + i = nalIdx // skip the matched start code (else a 4-byte code re-matches as 3-byte) + } + return false +} + // take in a segment and return a bunch of packets suitable for webrtc func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.PacketizedSegment, error) { @@ -116,6 +142,18 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti audioOutput := [][]byte{} // eosCh := make(chan struct{}) + // Slice-less video data (parameter sets / SEI metadata with no picture) + // held back for the next real frame. Streams with embedded closed captions + // carry a trailing caption SEI after some frames' slices; when h264parse + // re-parses the byte stream it splits that SEI into its own timestamp-less + // AU. Sent to WebRTC as a standalone "frame" it breaks strict decoders + // (iOS VideoToolbox errors on a picture-less access unit, and with PLI + // unanswered the picture stays broken until the next keyframe). Prepending + // it to the following frame is where a caption SEI normally lives, so + // caption-aware receivers still get it. A remainder at EOS has no frame to + // ride with and is dropped. + var pendingVideo []byte + videoappsink := app.SinkFromElement(videoSink) videoappsink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { @@ -131,15 +169,16 @@ func Packetize(ctx context.Context, cli *config.CLI, seg *bus.Seg) (*bus.Packeti samples := buffer.Bytes() - videoOutput = append(videoOutput, samples) + if !hasVideoSlice(samples) { + pendingVideo = append(pendingVideo, samples...) + return gst.FlowOK + } + if pendingVideo != nil { + samples = append(pendingVideo, samples...) + pendingVideo = nil + } - // clockTime := buffer.Duration() - // dur := clockTime.AsDuration() - // if dur != nil { - // log.Log(ctx, "video duration", "duration", *dur) - // } else { - // log.Error(ctx, "no video duration", "samples", len(samples)) - // } + videoOutput = append(videoOutput, samples) return gst.FlowOK }, diff --git a/pkg/media/packetize_test.go b/pkg/media/packetize_test.go index 586e79477..a96fe4ee7 100644 --- a/pkg/media/packetize_test.go +++ b/pkg/media/packetize_test.go @@ -1,17 +1,24 @@ package media import ( + "bytes" "context" + "encoding/binary" + "fmt" "io" "math/rand" "os" + "strings" "testing" "time" + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" "github.com/stretchr/testify/require" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/test/remote" ) @@ -71,6 +78,203 @@ func innerTestPacketize(t *testing.T, filename string, expectedVideo int, expect require.Equal(t, expectedDuration, packet.Duration) } +// captionSEINAL builds a minimal valid closed-caption SEI NAL (payload_type 4, +// user_data_registered_itu_t_t35, ATSC A/53 "GA94", two CEA-608 control-code +// pairs) — the kind of NAL a stream with embedded captions carries. Raw NAL +// bytes, no length/start-code framing. +func captionSEINAL() []byte { + t35 := []byte{ + 0xb5, // itu_t_t35_country_code: United States + 0x00, 0x31, // itu_t_t35_provider_code: ATSC + 'G', 'A', '9', '4', // user_identifier + 0x03, // user_data_type_code: cc_data + 0x40 | 0x02, // process_cc_data_flag, cc_count=2 + 0xff, // em_data + 0xfc, 0x94, 0xae, // cc_valid, NTSC field 1: ENM + 0xfc, 0x94, 0x20, // cc_valid, NTSC field 1: RCL + 0xff, // marker_bits + } + sei := []byte{0x06, 0x04, byte(len(t35))} + sei = append(sei, t35...) + return append(sei, 0x80) // rbsp_trailing_bits +} + +// makeTrailingCaptionSEIFlatMP4 synthesizes the sample shape MistServer +// produces for a stream with embedded closed captions: a flat fragmented MP4 +// whose video samples end with a caption SEI *after* the frame's slice. +// Returns the mp4 and how many samples carry the trailing SEI. When +// h264parse re-parses such a stream it splits each trailing SEI into its own +// slice-less, timestamp-less AU — the shape Packetize must not emit as a +// standalone video frame. +func makeTrailingCaptionSEIFlatMP4(t *testing.T, ctx context.Context, frames, seiEvery int) ([]byte, int) { + t.Helper() + gstinit.InitGST() + + // Encode raw AUs in avc stream-format (length-prefixed NALs), so appending + // a length-prefixed SEI to a sample is valid surgery. + type au struct { + data []byte + pts gst.ClockTime + dur gst.ClockTime + } + aus := []au{} + var vcaps *gst.Caps + encPipeline, err := gst.NewPipelineFromString(fmt.Sprintf( + "videotestsrc num-buffers=%d ! video/x-raw,width=320,height=240,framerate=30/1 ! x264enc tune=zerolatency speed-preset=ultrafast key-int-max=30 ! h264parse ! video/x-h264,stream-format=avc,alignment=au ! appsink name=sink", frames)) + require.NoError(t, err) + sinkEle, err := encPipeline.GetElementByName("sink") + require.NoError(t, err) + app.SinkFromElement(sinkEle).SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowEOS + } + if vcaps == nil { + vcaps = sample.GetCaps() + } + buffer := sample.GetBuffer() + aus = append(aus, au{ + data: append([]byte{}, buffer.Bytes()...), + pts: gst.ClockTime(buffer.PresentationTimestamp()), + dur: gst.ClockTime(buffer.Duration()), + }) + return gst.FlowOK + }, + }) + busErr := make(chan error, 1) + go func() { busErr <- HandleBusMessages(ctx, encPipeline) }() + require.NoError(t, encPipeline.SetState(gst.StatePlaying)) + require.NoError(t, <-busErr) + require.NoError(t, encPipeline.SetState(gst.StateNull)) + require.Len(t, aus, frames) + require.NotNil(t, vcaps) + + // x264enc offsets its output timestamps by a huge constant (its + // negative-DTS avoidance trick). Rebase to zero so the video timeline + // lines up with the audio track below — mismatched timelines make qtmux + // wait forever for the tracks to interleave. + base := aus[0].pts + for i := range aus { + aus[i].pts -= base + } + + // The surgery: append a caption SEI to every seiEvery-th sample. seiEvery + // must not divide frames evenly — a trailing SEI on the very last frame + // has no following frame to ride with and is (acceptably) dropped, which + // would confuse the survival assertion. + require.NotZero(t, frames%seiEvery) + sei := captionSEINAL() + prefixed := make([]byte, 4+len(sei)) + binary.BigEndian.PutUint32(prefixed, uint32(len(sei))) + copy(prefixed[4:], sei) + seiCount := 0 + for i := seiEvery - 1; i < len(aus)-1; i += seiEvery { + aus[i].data = append(aus[i].data, prefixed...) + seiCount++ + } + require.NotZero(t, seiCount) + + // Remux the doctored AUs (plus an Opus track, which Packetize requires) + // into a fragmented MP4. + audioBuffers := (frames*48000/30)/1024 + 1 + // Pad names are explicit: the appsrc's caps aren't known at parse time, so + // without them gst-parse guesses which mux request pad to link (and + // guesses wrong). + muxPipeline, err := gst.NewPipelineFromString(strings.Join([]string{ + "mp4mux name=mux fragment-duration=500 ! appsink name=sink", + "appsrc name=vsrc format=time ! mux.video_0", + fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! opusenc ! mux.audio_0", audioBuffers), + }, "\n")) + require.NoError(t, err) + vsrcEle, err := muxPipeline.GetElementByName("vsrc") + require.NoError(t, err) + require.NoError(t, vsrcEle.SetProperty("caps", vcaps)) + idx := 0 + app.SrcFromElement(vsrcEle).SetCallbacks(&app.SourceCallbacks{ + NeedDataFunc: func(self *app.Source, _ uint) { + if idx >= len(aus) { + self.EndStream() + return + } + a := aus[idx] + idx++ + buffer := gst.NewBufferFromBytes(a.data) + buffer.SetPresentationTimestamp(a.pts) + buffer.SetDuration(a.dur) + self.PushBuffer(buffer) + }, + }) + outSinkEle, err := muxPipeline.GetElementByName("sink") + require.NoError(t, err) + var out bytes.Buffer + app.SinkFromElement(outSinkEle).SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: WriterNewSample(ctx, &out), + }) + busErr2 := make(chan error, 1) + go func() { busErr2 <- HandleBusMessages(ctx, muxPipeline) }() + require.NoError(t, muxPipeline.SetState(gst.StatePlaying)) + require.NoError(t, <-busErr2) + require.NoError(t, muxPipeline.SetState(gst.StateNull)) + require.NotEmpty(t, out.Bytes()) + return out.Bytes(), seiCount +} + +// countCaptionSEIs counts caption SEI NALs (nal type 6, payload_type 4) in +// byte-stream H264 data. +func countCaptionSEIs(data []byte) int { + count := 0 + for i := 0; i+4 < len(data); i++ { + if data[i] != 0 || data[i+1] != 0 { + continue + } + var nalIdx int + if data[i+2] == 1 { + nalIdx = i + 3 + } else if data[i+2] == 0 && i+5 < len(data) && data[i+3] == 1 { + nalIdx = i + 4 + } else { + continue + } + if data[nalIdx]&0x1f == 6 && data[nalIdx+1] == 0x04 { + count++ + } + i = nalIdx // skip the matched start code (else a 4-byte code re-matches as 3-byte) + } + return count +} + +// TestPacketizeTrailingCaptionSEI: streams with embedded closed captions +// carry trailing caption SEIs that h264parse splits into slice-less, +// timestamp-less AUs. Packetize must fold those into the next real frame — +// not emit them as standalone video "frames", which inflate the frame count +// (skewing the sender's synthesized timing) and break strict WebRTC decoders +// (iOS VideoToolbox errors on a picture-less access unit). +func TestPacketizeTrailingCaptionSEI(t *testing.T) { + withNoGSTLeaks(t, func() { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + const frames = 60 + flat, seiCount := makeTrailingCaptionSEIFlatMP4(t, ctx, frames, 7) + + packet, err := Packetize(context.Background(), &config.CLI{}, &bus.Seg{Data: flat}) + require.NoError(t, err) + require.NotNil(t, packet) + + // Exactly one output sample per input frame — caption SEIs must not + // become frames of their own. + 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) + } + // ...and the captions must survive, riding with real frames. + require.Equal(t, seiCount, totalSEIs) + require.NotEmpty(t, packet.Audio) + }) +} + func TestPacketizeInvalid(t *testing.T) { // cur := goleak.IgnoreCurrent() // defer goleak.VerifyNone(t, cur) -- 2.51.2 From 1da503f397b80a1664b6043d05a70b317b3245ba Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 30 Jul 2026 20:10:56 -0700 Subject: [PATCH 2/5] 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 From 036f3b3a6b6d07821d6faa26e6b810be0ec85489 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 30 Jul 2026 20:23:12 -0700 Subject: [PATCH 3/5] webrtc: don't wedge Packetize on single-track segments MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ConcatDemuxBin pre-wires a video and an audio branch and waits for qtdemux to populate both — but a segment can legitimately have only one track: an encoder shedding everything but keyframes under bandwidth pressure produces GoPs whose audio never arrived. The demux only EOSes pads it actually created, so the trackless branch never saw data or EOS, the pipeline never completed, and Packetize burned its full 10s timeout — after which AddToWebRTC dropped the segment, leaving WebRTC viewers a multi-second hole (video AND audio) while HLS played on. On no-more-pads, complete any branch whose pad never appeared with an EOS so the pipeline drains promptly and the segment reaches viewers with the missing track empty. The handler uses the same swap-to-no-op indirection as pad-added so its pad references don't leak. Co-Authored-By: Claude Fable 5 --- pkg/media/concat_demux.go | 42 +++++++++++++++++++++++++++++++++++++ pkg/media/packetize_test.go | 35 +++++++++++++++++++++++++++++++ 2 files changed, 77 insertions(+) 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) -- 2.51.2 From 1d487f441838d1974c2eb5a00fd4c45e6389c2dd Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 31 Jul 2026 16:20:43 -0700 Subject: [PATCH 4/5] webrtc: address review feedback on duration + callback races MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A video-only segment reported Duration 0 — segDur is audio-derived — which threw off the sender's latency bookkeeping even though the segment occupies the sender for its full video span. Fall back to the video timeline when there's no audio. The demux signal handlers (pad-added and the new no-more-pads) were swapped for no-ops from the ctx goroutine with plain assignments the dispatch thread concurrently reads: a data race. Swap through atomic.Value instead. Notably NOT a mutex held across dispatch — that variant leaked GStreamer objects under the leak tracer (pads pinned alive along with their sticky tag events); an atomic store can't block behind a running handler, and dispatch behavior stays identical. Co-Authored-By: Claude Fable 5 --- pkg/media/concat_demux.go | 57 +++++++++++++++++++------------------ pkg/media/packetize.go | 14 +++++++-- pkg/media/packetize_test.go | 4 +++ 3 files changed, 46 insertions(+), 29 deletions(-) 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) { -- 2.51.2 From 0e74077e434ef147d48d7d4b665c5dbdc4f0f87f Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 31 Jul 2026 16:30:49 -0700 Subject: [PATCH 5/5] ci: retrigger after sharp/libvips download network flake linux-arm64 failed on `sharp: Installation error: read ECONNRESET` downloading libvips during pnpm install; siblings were fail-fast cancelled. No code change. Co-Authored-By: Claude Fable 5 -- 2.51.2