Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
15 kB · 396 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397package 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")
func TestPacketize(t *testing.T) { withNoGSTLeaks(t, func() { g, _ := errgroup.WithContext(context.Background()) for range streamplaceTestCount { g.Go(func() error { innerTestPacketize(t, getFixture("sample-segment.mp4"), 49, 40, time.Duration(800*time.Millisecond)) return nil }) } err := g.Wait() require.NoError(t, err) })}
func TestPacketizeMuxl(t *testing.T) { t.Run("BasicMuxl", func(t *testing.T) { withNoGSTLeaks(t, func() { filename := remote.RemoteFixture("c6b57a53fc5a2234dbdd388922f0e293d8063d2b30620321e974b7c85640f228/2026-03-17T19-02-08-607Z-muxl_segment_input.fmp4") innerTestPacketize(t, filename, 60, 50, time.Duration(1000*time.Millisecond)) }) }) t.Run("ThreeSecondSeg", func(t *testing.T) { withNoGSTLeaks(t, func() { filename2 := remote.RemoteFixture("5e2bca8cd42ad624d505c73f7b54ec761639416ef556d83145d37b730ed3606a/2026-04-11T22-24-11-527Z-packetize-input-019d7ea5-3d07-7606-b088-1a7bc315d009.mp4") innerTestPacketize(t, filename2, 180, 150, time.Duration(3000*time.Millisecond)) }) }) t.Run("TenSecondSeg", func(t *testing.T) { withNoGSTLeaks(t, func() { filename3 := remote.RemoteFixture("82d20ee62b02f1c3a727b3001f1fa939afb757f9f205fa438d7b5753e1253eef/2026-04-11T22-39-41-861Z-packetize-input-019d7eb3-6f24-776c-ba1b-2f909a2379d7.mp4") innerTestPacketize(t, filename3, 300, 502, time.Duration(10040*time.Millisecond)) }) })}
func innerTestPacketize(t *testing.T, filename string, expectedVideo int, expectedAudio int, expectedDuration time.Duration) { inputFile, err := os.Open(filename) require.NoError(t, err) defer inputFile.Close()
bs, err := io.ReadAll(inputFile) require.NoError(t, err)
testSeg := &bus.Seg{ Data: bs, Filepath: filename, }
packet, err := Packetize(context.Background(), &config.CLI{}, testSeg) require.NoError(t, err) require.NotNil(t, packet) require.Equal(t, expectedVideo, len(packet.Video)) require.Equal(t, expectedAudio, len(packet.Audio)) 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.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) require.NotEmpty(t, packet.Audio) })}
// 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))) })}
// 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) // 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) { 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) withNoGSTLeaks(t, func() { rng := rand.New(rand.NewSource(42)) randomData := make([]byte, 1024*1024) // 1MB _, err := rng.Read(randomData) require.NoError(t, err) packet, err := Packetize(context.Background(), &config.CLI{}, &bus.Seg{ Data: randomData, }) require.Error(t, err) require.Nil(t, packet) })}