diff --git a/pkg/config/config.go b/pkg/config/config.go index d850ef9b..74aef562 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -318,7 +318,7 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, &urfavecli.BoolFlag{ Name: "ll-hls", - Usage: "enable the experimental low-latency HLS origin for RTMP H.264/AAC streams", + Usage: "enable the experimental low-latency HLS origin for RTMP H.264/AAC and WHIP H.264/Opus streams", Value: false, Destination: &cli.LLHLS, Sources: urfavecli.EnvVars("SP_LL_HLS"), diff --git a/pkg/ingestframe/frame.go b/pkg/ingestframe/frame.go index d296115a..cab03864 100644 --- a/pkg/ingestframe/frame.go +++ b/pkg/ingestframe/frame.go @@ -169,20 +169,24 @@ func (fw *Writer) Manifest(payload []byte) error { return fw.WriteFrame(Manifest // timestamps use the worker's track timescale, keeping the wire protocol // independent of Go's duration representation. type LLFrame struct { - Presentation string `cbor:"presentation"` - Track string `cbor:"track"` - Generation uint64 `cbor:"generation"` - Timescale uint32 `cbor:"timescale,omitempty"` - MSN uint64 `cbor:"msn,omitempty"` - Part uint32 `cbor:"part,omitempty"` - Start uint64 `cbor:"start,omitempty"` - Duration uint64 `cbor:"duration,omitempty"` - Independent bool `cbor:"independent,omitempty"` - Data []byte `cbor:"data,omitempty"` + Presentation string `cbor:"presentation"` + Session uint64 `cbor:"session,omitempty"` + Track string `cbor:"track"` + Generation uint64 `cbor:"generation"` + Timescale uint32 `cbor:"timescale,omitempty"` + MSN uint64 `cbor:"msn,omitempty"` + Part uint32 `cbor:"part,omitempty"` + Start uint64 `cbor:"start,omitempty"` + Duration uint64 `cbor:"duration,omitempty"` + Independent bool `cbor:"independent,omitempty"` + ProgramDateTimeUnixNano int64 `cbor:"program_date_time_unix_nano,omitempty"` + FrameRate float64 `cbor:"frame_rate,omitempty"` + Channels int `cbor:"channels,omitempty"` + Data []byte `cbor:"data,omitempty"` } func (fw *Writer) writeLL(t Type, f LLFrame) error { - b, err := drisl.Marshal(f) + b, err := EncodeLLFrame(f) if err != nil { return fmt.Errorf("ingestframe: encode %s: %w", t, err) } @@ -194,6 +198,10 @@ func (fw *Writer) LLSegmentComplete(f LLFrame) error { return fw.writeLL(LLSegme func (fw *Writer) LLDiscontinuity(f LLFrame) error { return fw.writeLL(LLDiscontinuity, f) } func (fw *Writer) LLSessionEnd(f LLFrame) error { return fw.writeLL(LLSessionEnd, f) } +func EncodeLLFrame(f LLFrame) ([]byte, error) { + return drisl.Marshal(f) +} + func DecodeLLFrame(payload []byte) (LLFrame, error) { var f LLFrame if err := drisl.Unmarshal(payload, &f); err != nil { diff --git a/pkg/ingestframe/ll_test.go b/pkg/ingestframe/ll_test.go index 43a9318f..0600a060 100644 --- a/pkg/ingestframe/ll_test.go +++ b/pkg/ingestframe/ll_test.go @@ -8,7 +8,7 @@ import ( func TestLLFramesRoundTripMetadataAndBytes(t *testing.T) { var buf bytes.Buffer w := NewWriter(&buf) - if err := w.LLPart(LLFrame{Presentation: "p", Track: "v", Generation: 2, Timescale: 90000, MSN: 9, Part: 3, Start: 90, Duration: 30, Independent: true, Data: []byte("part")}); err != nil { + if err := w.LLPart(LLFrame{Presentation: "p", Session: 4, Track: "v", Generation: 2, Timescale: 90000, MSN: 9, Part: 3, Start: 90, Duration: 30, Independent: true, ProgramDateTimeUnixNano: 123456789, FrameRate: 120.5, Channels: 2, Data: []byte("part")}); err != nil { t.Fatal(err) } @@ -23,7 +23,7 @@ func TestLLFramesRoundTripMetadataAndBytes(t *testing.T) { if err != nil { t.Fatal(err) } - if got.Presentation != "p" || got.Track != "v" || got.Generation != 2 || got.Timescale != 90000 || got.MSN != 9 || got.Part != 3 || !got.Independent || !bytes.Equal(got.Data, []byte("part")) { + if got.Presentation != "p" || got.Session != 4 || got.Track != "v" || got.Generation != 2 || got.Timescale != 90000 || got.MSN != 9 || got.Part != 3 || !got.Independent || got.ProgramDateTimeUnixNano != 123456789 || got.FrameRate != 120.5 || got.Channels != 2 || !bytes.Equal(got.Data, []byte("part")) { t.Fatalf("frame = %+v", got) } } diff --git a/pkg/llhls/window.go b/pkg/llhls/window.go index 863bc46f..267559a5 100644 --- a/pkg/llhls/window.go +++ b/pkg/llhls/window.go @@ -45,11 +45,16 @@ const ( ) type Event struct { - Kind EventKind - Presentation string - Session uint64 - Track string - Generation uint64 + Kind EventKind + Presentation string + Session uint64 + Track string + Generation uint64 + // Bridge metadata is intentionally ignored by Window. Detached ingest uses it + // to carry CMAF timing and stream metadata across the worker boundary. + Timescale uint32 + FrameRate float64 + AudioChannels int MSN uint64 Part uint32 Start time.Duration @@ -127,6 +132,7 @@ type Window struct { partTarget time.Duration configuredTarget int64 configuredPartTarget time.Duration + dynamicTarget bool partHoldBack time.Duration configuredPartHoldBack time.Duration completionHold time.Duration @@ -207,6 +213,17 @@ func WithPlaylistDurations(parent, part time.Duration) Option { } } +// WithDynamicTargetDuration starts TARGETDURATION at the supplied minimum and +// raises it to cover the largest completed parent observed during the +// presentation. The value is rounded to the nearest whole second as required +// by HLS and never falls below one second. +func WithDynamicTargetDuration(minimum time.Duration) Option { + return func(w *Window) { + w.dynamicTarget = true + w.targetDuration = roundedDurationSeconds(minimum) + } +} + func NewWindow(opts ...Option) *Window { w := &Window{ tracks: make(map[string]*track), @@ -246,7 +263,7 @@ func (w *Window) Observe(ev Event) error { if ev.Kind != Init { return ErrStalePresentation } - if w.presentationSession != 0 && (ev.Session == 0 || ev.Session < w.presentationSession) { + if w.presentationSession != 0 && (ev.Session == 0 || ev.Session <= w.presentationSession) { return ErrStalePresentation } for _, t := range w.tracks { @@ -331,6 +348,12 @@ func (w *Window) Observe(ev Event) error { w.bytes += len(s.data) t.bytes += len(s.data) w.recordTrackBitrate(ev.Track, t, len(s.data), ev.Duration) + if w.dynamicTarget && ev.Duration > 0 { + observedTarget := roundedDurationSeconds(ev.Duration) + if observedTarget > w.targetDuration { + w.targetDuration = observedTarget + } + } if w.completionHold <= 0 { s.complete = true } else { diff --git a/pkg/llhls/window_test.go b/pkg/llhls/window_test.go index 25eac192..042760bc 100644 --- a/pkg/llhls/window_test.go +++ b/pkg/llhls/window_test.go @@ -645,6 +645,14 @@ func TestWindowRejectsStalePresentationSession(t *testing.T) { } } +func TestWindowRejectsEqualSessionStalePresentation(t *testing.T) { + w := NewWindow() + observeEvent(t, w, Event{Kind: Init, Presentation: "current", Session: 7, Track: "video", Generation: 1}) + if err := w.Observe(Event{Kind: Init, Presentation: "stale", Session: 7, Track: "video", Generation: 1}); err != ErrStalePresentation { + t.Fatalf("equal-session stale presentation error = %v", err) + } +} + func TestWindowResetsBitrateOnGenerationChange(t *testing.T) { w := NewWindow() observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "video", Generation: 1}) @@ -763,6 +771,53 @@ func TestPlaylistDurationsFreezeForPresentationAndRoundTargetDuration(t *testing } } +func TestDynamicPlaylistTargetDurationUsesObservedParentDurations(t *testing.T) { + w := NewWindow(WithDynamicTargetDuration(time.Second)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "video", Generation: 1}) + + playlist := w.Playlist("p", "video", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:1") { + t.Fatalf("initial target duration was not clamped to one second:\n%s", playlist) + } + + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Part: 0, Duration: 500 * time.Millisecond, Data: []byte("short")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Duration: 500 * time.Millisecond, Data: []byte("short")}) + playlist = w.Playlist("p", "video", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:1") { + t.Fatalf("short parent lowered the one-second floor:\n%s", playlist) + } + + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 2, Part: 0, Duration: 4 * time.Second, Data: []byte("long")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: 2, Duration: 4 * time.Second, Data: []byte("long")}) + playlist = w.Playlist("p", "video", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:4") { + t.Fatalf("target duration did not follow the four-second parent:\n%s", playlist) + } + + observeEvent(t, w, Event{Kind: Init, Presentation: "next", Track: "video", Generation: 1}) + playlist = w.Playlist("next", "video", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:1") { + t.Fatalf("target duration did not reset at the presentation boundary:\n%s", playlist) + } +} + +func TestDynamicPlaylistTargetDurationIsSharedAcrossTracks(t *testing.T) { + w := NewWindow(WithDynamicTargetDuration(time.Second)) + for _, track := range []string{"video", "audio"} { + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: track, Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: track, Generation: 1, MSN: 1, Part: 0, Duration: time.Second, Data: []byte(track)}) + } + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Duration: time.Second, Data: []byte("video")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Duration: 4 * time.Second, Data: []byte("audio")}) + + for _, track := range []string{"video", "audio"} { + playlist := w.Playlist("p", track, completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:4") { + t.Errorf("%s playlist did not use the window-wide observed target:\n%s", track, playlist) + } + } +} + func TestPlaylistUsesConfiguredPartTargetWithoutGrowing(t *testing.T) { w := NewWindow() observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) diff --git a/pkg/media/cmaf.go b/pkg/media/cmaf.go index 159891c5..d2f4655d 100644 --- a/pkg/media/cmaf.go +++ b/pkg/media/cmaf.go @@ -30,6 +30,7 @@ type cmafTrackSink struct { session uint64 track string window *llhls.Window + publish func(llhls.Event) error generation uint64 nextMSN uint64 initialized bool @@ -50,6 +51,8 @@ type cmafTrackSink struct { // times provide the offset for each parent. programDateTimeBase time.Time timescale uint32 + videoFrameRate float64 + audioChannels int } type cmafPendingPart struct { @@ -112,7 +115,8 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { } else { s.timescale = timescale } - if err := s.window.Observe(llhls.Event{ + s.audioChannels = audioChannels + if err := s.observe(llhls.Event{ Kind: llhls.Init, Presentation: s.presentation, Session: s.session, @@ -123,7 +127,7 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { return fmt.Errorf("publish CMAF init: %w", err) } if audioChannels > 0 { - s.window.SetAudioConfig(llhls.AudioConfig{Channels: audioChannels}) + s.setAudioConfig(llhls.AudioConfig{Channels: audioChannels}) } s.initialized = true } @@ -176,7 +180,11 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { if s.track == "video" { for _, timing := range timings { if s.videoTrackIDs[timing.TrackID] { - s.window.SetVideoFrameRate(cmafVideoFrameRate(timing, s.timescale)) + rate := cmafVideoFrameRate(timing, s.timescale) + if rate > s.videoFrameRate { + s.videoFrameRate = rate + } + s.setVideoFrameRate(rate) break } } @@ -302,7 +310,7 @@ func (s *cmafTrackSink) publishPart(part cmafPendingPart) error { if s.partTarget > 0 && part.duration > s.partTarget { return fmt.Errorf("CMAF part duration %s exceeds PART-TARGET %s", part.duration, s.partTarget) } - if err := s.window.Observe(llhls.Event{ + if err := s.observe(llhls.Event{ Kind: llhls.Part, Presentation: s.presentation, Session: s.session, @@ -330,7 +338,7 @@ func (s *cmafTrackSink) completeParent() error { return err } s.pendingPart = cmafPendingPart{} - if err := s.window.Observe(llhls.Event{ + if err := s.observe(llhls.Event{ Kind: llhls.SegmentComplete, Presentation: s.presentation, Session: s.session, @@ -352,6 +360,31 @@ func (s *cmafTrackSink) completeParent() error { return nil } +func (s *cmafTrackSink) observe(ev llhls.Event) error { + ev.Timescale = s.timescale + ev.FrameRate = s.videoFrameRate + ev.AudioChannels = s.audioChannels + if s.publish != nil { + return s.publish(ev) + } + if s.window == nil { + return errors.New("CMAF event sink has no destination") + } + return s.window.Observe(ev) +} + +func (s *cmafTrackSink) setVideoFrameRate(fps float64) { + if s.window != nil { + s.window.SetVideoFrameRate(fps) + } +} + +func (s *cmafTrackSink) setAudioConfig(config llhls.AudioConfig) { + if s.window != nil { + s.window.SetAudioConfig(config) + } +} + func isCMAFInit(data []byte) bool { return len(data) >= 8 && string(data[4:8]) == "ftyp" } diff --git a/pkg/media/cmaf_test.go b/pkg/media/cmaf_test.go index 672edd35..4e8fb6fd 100644 --- a/pkg/media/cmaf_test.go +++ b/pkg/media/cmaf_test.go @@ -58,6 +58,36 @@ func TestCMAFPartPublishesFullSizedFirstPartImmediately(t *testing.T) { } } +func TestCMAFSinkPublishesWorkerMetadata(t *testing.T) { + var got llhls.Event + state := &cmafTrackSink{ + presentation: "test", + session: 4, + track: "audio", + generation: 2, + timescale: 48000, + videoFrameRate: 59.94, + audioChannels: 2, + publish: func(ev llhls.Event) error { + got = ev + return nil + }, + } + when := time.Unix(1700000000, 123) + if err := state.publishPart(cmafPendingPart{ + data: []byte("part"), + start: time.Second, + duration: 20 * time.Millisecond, + programDateTime: when, + set: true, + }); err != nil { + t.Fatal(err) + } + if got.Session != 4 || got.Timescale != 48000 || got.FrameRate != 59.94 || got.AudioChannels != 2 || !got.ProgramDateTime.Equal(when) { + t.Fatalf("worker event metadata = %+v", got) + } +} + func TestCMAFMuxEmitsFragmentedBufferLists(t *testing.T) { gstinit.InitGST() if gst.Find("cmafmux") == nil || gst.Find("x264enc") == nil { diff --git a/pkg/media/frame_server.go b/pkg/media/frame_server.go index a60d7e38..48cea19b 100644 --- a/pkg/media/frame_server.go +++ b/pkg/media/frame_server.go @@ -14,11 +14,13 @@ import ( "stream.place/streamplace/pkg/log" ) -// workerFrameBuffer bounds how many signed segments a worker holds while main is -// disconnected — ~10 min at one segment per ~1s GoP. Beyond this the oldest are -// dropped (a main outage longer than this loses the oldest tail, loudly). +// workerFrameBuffer bounds how many frames a worker holds while main is +// disconnected. LL-HLS parts and metadata consume entries alongside signed +// segments, so this is a frame bound rather than a fixed time guarantee. const workerFrameBuffer = 600 +const workerFrameWriteTimeout = 100 * time.Millisecond + // workerDrainGrace bounds how long a worker lingers after its stream ends waiting // for main to drain the buffer. Generous enough for a main restart/upgrade; an // orphaned worker (main never returns) exits after it. @@ -53,6 +55,8 @@ type bufferedFrame struct { // and re-buffers that frame, so a hard main disconnect degrades to buffering. type frameServer struct { mu sync.Mutex + llInits []bufferedFrame + llInitTracks map[string]int pending []bufferedFrame conn net.Conn w *ingestframe.Writer @@ -64,31 +68,93 @@ type frameServer struct { // newFrameServer creates a server that buffers up to maxBuf frames while no // client is attached. func newFrameServer(maxBuf int) *frameServer { - return &frameServer{maxBuf: maxBuf} + return &frameServer{maxBuf: maxBuf, llInitTracks: make(map[string]int)} } func (s *frameServer) push(typ ingestframe.Type, payload []byte) { s.mu.Lock() defer s.mu.Unlock() + rememberedInit := false + if typ == ingestframe.LLInit { + if frame, err := ingestframe.DecodeLLFrame(payload); err == nil && frame.Track != "" { + if s.llInitTracks == nil { + s.llInitTracks = make(map[string]int) + } + buffered := bufferedFrame{typ, bytes.Clone(payload)} + if index, ok := s.llInitTracks[frame.Track]; ok { + s.llInits[index] = buffered + } else { + s.llInitTracks[frame.Track] = len(s.llInits) + s.llInits = append(s.llInits, buffered) + } + rememberedInit = true + } + } if s.w != nil { - if err := s.w.WriteFrame(typ, payload); err == nil { + if err := writeWorkerFrame(s.conn, s.w, typ, payload); err == nil { return } // Client gone; drop it and buffer this frame instead. + conn := s.conn s.conn, s.w = nil, nil + if conn != nil { + _ = conn.Close() + } + } + if typ != ingestframe.LLInit || !rememberedInit { + s.pending = append(s.pending, bufferedFrame{typ, bytes.Clone(payload)}) } - s.pending = append(s.pending, bufferedFrame{typ, bytes.Clone(payload)}) for len(s.pending) > s.maxBuf { s.pending = s.pending[1:] s.dropped++ } } +func writeWorkerFrame(conn net.Conn, writer *ingestframe.Writer, typ ingestframe.Type, payload []byte) error { + if conn == nil || writer == nil { + return fmt.Errorf("missing frame connection") + } + if err := conn.SetWriteDeadline(time.Now().Add(workerFrameWriteTimeout)); err != nil { + return err + } + defer conn.SetWriteDeadline(time.Time{}) + return writer.WriteFrame(typ, payload) +} + func (s *frameServer) Segment(seg []byte) error { s.push(ingestframe.Segment, seg); return nil } func (s *frameServer) End() error { s.push(ingestframe.End, nil); return nil } func (s *frameServer) Error(msg string) error { s.push(ingestframe.Error, []byte(msg)); return nil } func (s *frameServer) Answer(sdp string) error { s.push(ingestframe.Answer, []byte(sdp)); return nil } +func (s *frameServer) writeLL(typ ingestframe.Type, frame ingestframe.LLFrame) error { + payload, err := ingestframe.EncodeLLFrame(frame) + if err != nil { + return err + } + s.push(typ, payload) + return nil +} + +func (s *frameServer) LLInit(frame ingestframe.LLFrame) error { + return s.writeLL(ingestframe.LLInit, frame) +} + +func (s *frameServer) LLPart(frame ingestframe.LLFrame) error { + return s.writeLL(ingestframe.LLPart, frame) +} + +func (s *frameServer) LLSegmentComplete(frame ingestframe.LLFrame) error { + return s.writeLL(ingestframe.LLSegmentComplete, frame) +} + +func (s *frameServer) LLDiscontinuity(frame ingestframe.LLFrame) error { + return s.writeLL(ingestframe.LLDiscontinuity, frame) +} + +func (s *frameServer) LLSessionEnd(frame ingestframe.LLFrame) error { + return s.writeLL(ingestframe.LLSessionEnd, frame) +} + // dropped reports how many buffered frames were discarded because the buffer // overflowed (main was disconnected longer than the buffer window). func (s *frameServer) droppedCount() int { @@ -104,8 +170,15 @@ func (s *frameServer) attach(conn net.Conn) { s.mu.Lock() defer s.mu.Unlock() w := ingestframe.NewWriter(conn) + for _, f := range s.llInits { + if err := writeWorkerFrame(conn, w, f.typ, f.payload); err != nil { + _ = conn.Close() + return + } + } for _, f := range s.pending { - if err := w.WriteFrame(f.typ, f.payload); err != nil { + if err := writeWorkerFrame(conn, w, f.typ, f.payload); err != nil { + _ = conn.Close() return } } @@ -135,7 +208,10 @@ func (s *frameServer) waitDrained(ctx context.Context, grace time.Duration) { case <-ctx.Done(): return case <-deadline.C: - log.Warn(ctx, "ingest worker: main never drained the frame buffer; exiting", "pending", len(s.pending)) + s.mu.Lock() + pending := len(s.pending) + s.mu.Unlock() + log.Warn(ctx, "ingest worker: main never drained the frame buffer; exiting", "pending", pending) return case <-tick.C: } diff --git a/pkg/media/frame_server_test.go b/pkg/media/frame_server_test.go index 22810e2d..dd98dd9f 100644 --- a/pkg/media/frame_server_test.go +++ b/pkg/media/frame_server_test.go @@ -4,8 +4,10 @@ import ( "context" "fmt" "net" + "os" "path/filepath" "testing" + "time" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/ingestframe" @@ -16,7 +18,10 @@ import ( // the frameServer can flush its buffer without a concurrent drainer. func unixPair(t *testing.T) (client net.Conn, server net.Conn) { t.Helper() - sock := filepath.Join(t.TempDir(), "p.sock") + tempDir, err := os.MkdirTemp("/tmp", "sp-") + require.NoError(t, err) + t.Cleanup(func() { os.RemoveAll(tempDir) }) + sock := filepath.Join(tempDir, "p.sock") ln, err := net.Listen("unix", sock) require.NoError(t, err) defer ln.Close() @@ -95,6 +100,69 @@ func TestFrameServerBufferFlushReconnect(t *testing.T) { require.Equal(t, 0, srv.droppedCount(), "ample buffer drops nothing across a brief restart") } +func TestFrameServerReplaysLatestLLInitOnReconnect(t *testing.T) { + srv := newFrameServer(1000) + initFrame := ingestframe.LLFrame{ + Presentation: "whip-1-test", + Track: "video", + Generation: 1, + Data: []byte("init"), + } + partFrame := initFrame + partFrame.MSN = 1 + partFrame.Data = []byte("part") + initPayload, err := ingestframe.EncodeLLFrame(initFrame) + require.NoError(t, err) + partPayload, err := ingestframe.EncodeLLFrame(partFrame) + require.NoError(t, err) + + // The first init is delivered before main disconnects. The next part is + // buffered, so a fresh main must receive the init again before that part. + clientA, serverA := unixPair(t) + srv.attach(serverA) + srv.push(ingestframe.LLInit, initPayload) + _, _, err = ingestframe.NewReader(clientA).ReadFrame() + require.NoError(t, err) + srv.detachConn(serverA) + clientA.Close() + serverA.Close() + srv.push(ingestframe.LLPart, partPayload) + + clientB, serverB := unixPair(t) + srv.attach(serverB) + reader := ingestframe.NewReader(clientB) + typ, payload, err := reader.ReadFrame() + require.NoError(t, err) + require.Equal(t, ingestframe.LLInit, typ) + require.Equal(t, initPayload, payload) + typ, payload, err = reader.ReadFrame() + require.NoError(t, err) + require.Equal(t, ingestframe.LLPart, typ) + require.Equal(t, partPayload, payload) +} + +func TestFrameServerBoundsBlockedClientWrite(t *testing.T) { + client, server := net.Pipe() + defer client.Close() + defer server.Close() + + srv := newFrameServer(3) + srv.mu.Lock() + srv.conn = server + srv.w = ingestframe.NewWriter(server) + srv.mu.Unlock() + + start := time.Now() + require.NoError(t, srv.Segment(seg(0))) + require.Less(t, time.Since(start), 500*time.Millisecond, "blocked client write exceeded its deadline") + + srv.mu.Lock() + defer srv.mu.Unlock() + require.Nil(t, srv.conn) + require.Nil(t, srv.w) + require.Len(t, srv.pending, 1) +} + // TestFrameServerDropsOldestBeyondBound: a main outage longer than the buffer // window drops the OLDEST frames (bounded memory), loudly via droppedCount — // never grows without limit. diff --git a/pkg/media/ingest_daemon.go b/pkg/media/ingest_daemon.go index cab8825b..f358f76b 100644 --- a/pkg/media/ingest_daemon.go +++ b/pkg/media/ingest_daemon.go @@ -15,6 +15,7 @@ import ( "github.com/google/uuid" "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/llhls" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/spmetrics" ) @@ -84,7 +85,13 @@ func SpawnIngestWorkerDetached(cfg IngestWorkerConfig, media *os.File) (*os.Proc // hold this worker's process handle, kill it by PID if needed. Best-effort: // without it the worker still runs, just without restart-surviving ban // enforcement. - if err := writeWorkerMeta(cfg.SocketPath, workerMeta{StreamerDID: cfg.StreamerDID, PID: cmd.Process.Pid}); err != nil { + if err := writeWorkerMeta(cfg.SocketPath, workerMeta{ + StreamerDID: cfg.StreamerDID, + PID: cmd.Process.Pid, + LLHLS: cfg.LLHLS, + LLHLSPresentation: cfg.LLHLSPresentation, + LLHLSSession: cfg.LLHLSSession, + }); err != nil { log.Warn(context.Background(), "ingest worker: write resume metadata failed", "socket", cfg.SocketPath, "error", err) } return cmd.Process, nil @@ -93,8 +100,11 @@ func SpawnIngestWorkerDetached(cfg IngestWorkerConfig, media *os.File) (*os.Proc // workerMeta is the resume sidecar written next to a detached worker's socket: // enough for a restarting main to enforce bans on a worker it didn't spawn. type workerMeta struct { - StreamerDID string `json:"streamer_did"` - PID int `json:"pid"` + StreamerDID string `json:"streamer_did"` + PID int `json:"pid"` + LLHLS bool `json:"ll_hls,omitempty"` + LLHLSPresentation string `json:"ll_hls_presentation,omitempty"` + LLHLSSession uint64 `json:"ll_hls_session,omitempty"` } func workerMetaPath(socketPath string) string { @@ -367,8 +377,24 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, if err != nil { return "", err } + var llWindow *llhls.Window + llPresentation := "" + if webRTCLLHLSAvailable(mm.cli) { + llSession := mm.nextIngestSession() + llPresentation = fmt.Sprintf("whip-%d-%s", llSession, uuid.NewString()) + llWindow = mm.replaceLLWindow(ms.Streamer()) + cfg.LLHLS = true + cfg.LLHLSPresentation = llPresentation + cfg.LLHLSSession = llSession + } + cleanupLLHLS := func() { + if llWindow != nil { + mm.removeLLWindow(ms.Streamer(), llPresentation, llWindow) + } + } dir, err := mm.ingestWorkerSocketDir() if err != nil { + cleanupLLHLS() return "", err } cfg.SocketPath = filepath.Join(dir, uuid.NewString()+".sock") @@ -377,6 +403,7 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, proc, err := SpawnIngestWorkerDetached(cfg, nil) // worker owns the PeerConnection if err != nil { + cleanupLLHLS() return "", fmt.Errorf("spawn detached whip worker: %w", err) } spmetrics.IngestWorkerStarts.WithLabelValues("whip").Inc() @@ -388,6 +415,7 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, conn, err := dialWorkerSocket(answerCtx, cfg.SocketPath) if err != nil { _ = proc.Kill() + cleanupLLHLS() return "", fmt.Errorf("connect to whip worker: %w", err) } // One Reader owns this connection for its whole lifetime: the streaming CBOR @@ -401,13 +429,15 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, if err != nil { conn.Close() _ = proc.Kill() + cleanupLLHLS() return "", err } _ = conn.SetReadDeadline(time.Time{}) // clear; streaming has no deadline - // Consume the signed segments in the background; the HTTP handler returns the - // answer now and the WebRTC media establishes directly to the worker. + // Consume signed segments and LL-HLS events in the background; the HTTP handler + // returns the answer now and the WebRTC media establishes directly to the worker. go func() { + defer cleanupLLHLS() // Ban / key revocation: watch on the detached worker's behalf and kill it. // Scoped to this consume's lifetime so it doesn't outlive the stream. wctx, wcancel := context.WithCancel(ctx) @@ -459,9 +489,20 @@ func (mm *MediaManager) ResumeDetachedWorkers(ctx context.Context) { streamer = meta.StreamerDID } log.Log(ctx, "resuming detached ingest worker", "socket", sock, "streamer", streamer) + var llWindow *llhls.Window + llPresentation := "" + if merr == nil && meta.LLHLS && meta.LLHLSPresentation != "" && meta.StreamerDID != "" { + llPresentation = meta.LLHLSPresentation + llWindow = mm.replaceLLWindow(meta.StreamerDID) + } go func() { wctx, wcancel := context.WithCancel(ctx) defer wcancel() + defer func() { + if llWindow != nil { + mm.removeLLWindow(streamer, llPresentation, llWindow) + } + }() // Re-arm ban enforcement for a worker we didn't spawn. We have no process // handle, so kill by the PID in the sidecar (guarded against PID reuse). // Without metadata we can still drain the worker, just not enforce bans. diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 8fd05e10..e52341ac 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -215,6 +215,10 @@ func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, fr *ingestframe log.Error(ctx, "ingest worker: segment handler failed", "streamer", streamer, "error", serr) } } + case ingestframe.LLInit, ingestframe.LLPart, ingestframe.LLSegmentComplete, ingestframe.LLDiscontinuity, ingestframe.LLSessionEnd: + if err := mm.observeWorkerLLHLSFrame(streamer, typ, payload); err != nil { + log.Error(ctx, "ingest worker: LL-HLS frame failed", "streamer", streamer, "type", typ, "error", err) + } case ingestframe.End: sawEnd = true case ingestframe.Error: diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index d22ef023..d427411e 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -56,11 +56,9 @@ type IngestWorkerConfig struct { // forwarded verbatim, no reconstruction. KeyPEM []byte `json:"key_pem"` CertPEM []byte `json:"cert_pem"` - // Manifest is the C2PA manifest JSON, built ONCE by main at stream start. - // muxl-sign stamps each segment's signing time into it as it signs. NOTE: - // static for the worker's lifetime — mid-stream manifest changes (e.g. a - // pre-live → live transition) don't yet cross the boundary; that needs a - // control channel and is tracked as future work. + // Manifest is the initial C2PA manifest JSON. Main can replace it over the + // worker socket while the stream is running; muxl-sign reads the current + // value for each segment and stamps that segment's signing time into it. Manifest []byte `json:"manifest"` // Node transcode signer + broadcaster identity. When set, the worker completes @@ -112,6 +110,12 @@ type IngestWorkerConfig struct { // generates the answer and emits it as the first frame (ingestframe.Answer) // so main can return it to the client before consuming segments. OfferSDP string `json:"offer_sdp,omitempty"` + // LLHLS asks a WHIP worker to emit raw LL-HLS events alongside signed + // segments. The presentation identity is allocated by main so the main + // process can create the playback window before the first event arrives. + LLHLS bool `json:"ll_hls,omitempty"` + LLHLSPresentation string `json:"ll_hls_presentation,omitempty"` + LLHLSSession uint64 `json:"ll_hls_session,omitempty"` } // IngestTransportWHIP is the cfg.Transport value selecting the WHIP worker. @@ -122,7 +126,7 @@ const IngestTransportWHIP = "whip" // handed its config over the handshake, else local disk under DataDir). Shared // by the MP4 and WHIP workers so both record to the same place main would. func (cfg IngestWorkerConfig) workerCLI() *config.CLI { - cli := &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir} + cli := &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir, LLHLS: cfg.LLHLS} if cfg.S3 != nil { cli.SetS3Config(*cfg.S3) } diff --git a/pkg/media/live_window.go b/pkg/media/live_window.go index aeb493bb..d22ea47c 100644 --- a/pkg/media/live_window.go +++ b/pkg/media/live_window.go @@ -35,7 +35,8 @@ func newLLWindow() *llhls.Window { return llhls.NewWindow( llhls.WithMaxSegments(llhlsWindowSegments), llhls.WithMaxBytes(llhlsWindowBytes), - llhls.WithPlaylistDurations(llhlsParentDuration, llhlsPartTarget), + llhls.WithDynamicTargetDuration(time.Second), + llhls.WithPlaylistDurations(0, llhlsPartTarget), llhls.WithPartHoldBack(llhlsLivePartHoldBack), llhls.WithSegmentCompletionDelay(llhlsCompletionHold), ) @@ -44,6 +45,9 @@ func newLLWindow() *llhls.Window { func (mm *MediaManager) replaceLLWindow(did string) *llhls.Window { mm.llWindowsMut.Lock() defer mm.llWindowsMut.Unlock() + if mm.llWindows == nil { + mm.llWindows = make(map[string]*llhls.Window) + } w := newLLWindow() mm.llWindows[did] = w return w @@ -62,7 +66,7 @@ func (mm *MediaManager) removeLLWindow(did, presentation string, expected *llhls mm.llWindowsMut.Lock() defer mm.llWindowsMut.Unlock() window := mm.llWindows[did] - if window == expected && window != nil && window.Presentation() == presentation { + if window == expected && window != nil && (window.Presentation() == "" || window.Presentation() == presentation) { delete(mm.llWindows, did) } } diff --git a/pkg/media/live_window_ll_test.go b/pkg/media/live_window_ll_test.go index 36c0e820..7d63ad4b 100644 --- a/pkg/media/live_window_ll_test.go +++ b/pkg/media/live_window_ll_test.go @@ -32,9 +32,9 @@ func TestLLWindowUsesMobilePlaybackHoldBack(t *testing.T) { })) playlist := w.Playlist("p", "video", func(uint64, uint32) string { return "part.m4s" }, func(uint64) string { return "segment.m4s" }, "init.mp4", nil) - require.Contains(t, playlist, "#EXT-X-TARGETDURATION:2") + require.Contains(t, playlist, "#EXT-X-TARGETDURATION:1") require.Contains(t, playlist, "PART-HOLD-BACK=5.500000") - require.Contains(t, playlist, "HOLD-BACK=6.000000") + require.Contains(t, playlist, "HOLD-BACK=3.000000") } func TestRemoveLLWindowOnlyRemovesMatchingPresentation(t *testing.T) { diff --git a/pkg/media/llhls_ingest.go b/pkg/media/llhls_ingest.go new file mode 100644 index 00000000..7334edf3 --- /dev/null +++ b/pkg/media/llhls_ingest.go @@ -0,0 +1,247 @@ +package media + +import ( + "fmt" + "sync" + "time" + + "github.com/go-gst/go-gst/gst" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/llhls" +) + +// llhlsIngestOutput is the boundary between an ingest pipeline and its LL-HLS +// destination. The in-process path observes directly into a Window; an isolated +// worker serializes the same events onto its reconnectable frame stream for main +// to observe. +type llhlsIngestOutput struct { + presentation string + session uint64 + window *llhls.Window + publish func(llhls.Event) error + done func() +} + +func webRTCLLHLSAvailable(cli *config.CLI) bool { + if cli == nil || !cli.LLHLS { + return false + } + for _, name := range []string{"isofmp4mux", "opusdec", "fdkaacenc"} { + if gst.Find(name) == nil { + return false + } + } + return true +} + +type llhlsFrameWriter interface { + LLInit(ingestframe.LLFrame) error + LLPart(ingestframe.LLFrame) error + LLSegmentComplete(ingestframe.LLFrame) error + LLDiscontinuity(ingestframe.LLFrame) error + LLSessionEnd(ingestframe.LLFrame) error +} + +type workerLLHLSOutput struct { + writer llhlsFrameWriter + mu sync.Mutex + tracks map[string]uint64 + presentation string + session uint64 +} + +func newWorkerLLHLSOutput(writer llhlsFrameWriter, presentation string, session uint64) *workerLLHLSOutput { + return &workerLLHLSOutput{ + writer: writer, + tracks: make(map[string]uint64), + presentation: presentation, + session: session, + } +} + +func (o *workerLLHLSOutput) publish(ev llhls.Event) error { + frame := llhlsEventToFrame(ev) + o.mu.Lock() + if ev.Track != "" && ev.Kind == llhls.Init { + o.tracks[ev.Track] = ev.Generation + } + o.mu.Unlock() + + switch ev.Kind { + case llhls.Init: + return o.writer.LLInit(frame) + case llhls.Part: + return o.writer.LLPart(frame) + case llhls.SegmentComplete: + return o.writer.LLSegmentComplete(frame) + case llhls.Discontinuity: + return o.writer.LLDiscontinuity(frame) + case llhls.SessionEnd: + return o.writer.LLSessionEnd(frame) + default: + return fmt.Errorf("unsupported LL-HLS event kind %d", ev.Kind) + } +} + +func (o *workerLLHLSOutput) done() { + o.mu.Lock() + track, generation := "video", o.tracks["video"] + if generation == 0 { + for candidate, candidateGeneration := range o.tracks { + track, generation = candidate, candidateGeneration + break + } + } + presentation, session := o.presentation, o.session + o.mu.Unlock() + if generation == 0 { + return + } + if err := o.publish(llhls.Event{ + Kind: llhls.SessionEnd, + Presentation: presentation, + Session: session, + Track: track, + Generation: generation, + }); err != nil { + return + } +} + +func llhlsEventToFrame(ev llhls.Event) ingestframe.LLFrame { + frame := ingestframe.LLFrame{ + Presentation: ev.Presentation, + Session: ev.Session, + Track: ev.Track, + Generation: ev.Generation, + Timescale: ev.Timescale, + MSN: ev.MSN, + Part: ev.Part, + Start: durationToLLTicks(ev.Start, ev.Timescale), + Duration: durationToLLTicks(ev.Duration, ev.Timescale), + Independent: ev.Independent, + FrameRate: ev.FrameRate, + Channels: ev.AudioChannels, + Data: ev.Data, + } + if !ev.ProgramDateTime.IsZero() { + frame.ProgramDateTimeUnixNano = ev.ProgramDateTime.UnixNano() + } + return frame +} + +func llhlsFrameToEvent(kind ingestframe.Type, frame ingestframe.LLFrame) (llhls.Event, error) { + var eventKind llhls.EventKind + switch kind { + case ingestframe.LLInit: + eventKind = llhls.Init + case ingestframe.LLPart: + eventKind = llhls.Part + case ingestframe.LLSegmentComplete: + eventKind = llhls.SegmentComplete + case ingestframe.LLDiscontinuity: + eventKind = llhls.Discontinuity + case ingestframe.LLSessionEnd: + eventKind = llhls.SessionEnd + default: + return llhls.Event{}, fmt.Errorf("unsupported LL-HLS worker frame %s", kind) + } + start, err := durationFromLLTicks(frame.Start, frame.Timescale) + if err != nil { + return llhls.Event{}, fmt.Errorf("invalid LL-HLS worker start: %w", err) + } + duration, err := durationFromLLTicks(frame.Duration, frame.Timescale) + if err != nil { + return llhls.Event{}, fmt.Errorf("invalid LL-HLS worker duration: %w", err) + } + if frame.Presentation == "" || frame.Track == "" || frame.Generation == 0 { + return llhls.Event{}, fmt.Errorf("invalid LL-HLS worker frame: missing presentation, track, or generation") + } + event := llhls.Event{ + Kind: eventKind, + Presentation: frame.Presentation, + Session: frame.Session, + Track: frame.Track, + Generation: frame.Generation, + Timescale: frame.Timescale, + MSN: frame.MSN, + Part: frame.Part, + Start: start, + Duration: duration, + Independent: frame.Independent, + FrameRate: frame.FrameRate, + AudioChannels: frame.Channels, + Data: frame.Data, + } + if frame.ProgramDateTimeUnixNano != 0 { + event.ProgramDateTime = time.Unix(0, frame.ProgramDateTimeUnixNano) + } + return event, nil +} + +func durationToLLTicks(value time.Duration, timescale uint32) uint64 { + if value <= 0 { + return 0 + } + if timescale == 0 { + return uint64(value) + } + return uint64(value) * uint64(timescale) / uint64(time.Second) +} + +func durationFromLLTicks(value uint64, timescale uint32) (time.Duration, error) { + if value == 0 { + return 0, nil + } + if timescale == 0 { + if value > maxDurationNanos { + return 0, fmt.Errorf("tick value %d exceeds duration range", value) + } + return time.Duration(value), nil + } + scale := uint64(timescale) + wholeSeconds := value / scale + if wholeSeconds > maxDurationNanos/uint64(time.Second) { + return 0, fmt.Errorf("tick value %d exceeds duration range", value) + } + nanos := wholeSeconds * uint64(time.Second) + fractionalNanos := (value % scale) * uint64(time.Second) / scale + if nanos > maxDurationNanos-fractionalNanos { + return 0, fmt.Errorf("tick value %d exceeds duration range", value) + } + return time.Duration(nanos + fractionalNanos), nil +} + +const maxDurationNanos = uint64(1<<63 - 1) + +func (mm *MediaManager) observeWorkerLLHLSFrame(streamer string, kind ingestframe.Type, payload []byte) error { + frame, err := ingestframe.DecodeLLFrame(payload) + if err != nil { + return err + } + event, err := llhlsFrameToEvent(kind, frame) + if err != nil { + return err + } + window := mm.GetLLWindow(streamer) + if window == nil { + if event.Kind != llhls.Init { + return fmt.Errorf("LL-HLS worker event arrived before init") + } + window = mm.replaceLLWindow(streamer) + } + if event.FrameRate > 0 { + window.SetVideoFrameRate(event.FrameRate) + } + if event.AudioChannels > 0 { + window.SetAudioConfig(llhls.AudioConfig{Channels: event.AudioChannels}) + } + if err := window.Observe(event); err != nil { + return err + } + if event.Kind == llhls.SessionEnd { + mm.removeLLWindow(streamer, event.Presentation, window) + } + return nil +} diff --git a/pkg/media/llhls_ingest_test.go b/pkg/media/llhls_ingest_test.go new file mode 100644 index 00000000..a2dbcbfa --- /dev/null +++ b/pkg/media/llhls_ingest_test.go @@ -0,0 +1,136 @@ +package media + +import ( + "testing" + "time" + + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/llhls" +) + +func TestObserveWorkerLLHLSFramePopulatesWindow(t *testing.T) { + const streamer = "did:key:z6MkWorkerLLHLSWindowTest" + mm := &MediaManager{} + window := mm.replaceLLWindow(streamer) + + frames := []struct { + typ ingestframe.Type + event llhls.Event + }{ + {typ: ingestframe.LLInit, event: llhls.Event{ + Kind: llhls.Init, + Presentation: "whip-1-test", + Session: 1, + Track: "video", + Generation: 1, + Timescale: 90000, + FrameRate: 120, + Data: []byte("video-init"), + }}, + {typ: ingestframe.LLInit, event: llhls.Event{ + Kind: llhls.Init, + Presentation: "whip-1-test", + Session: 1, + Track: "audio", + Generation: 1, + Timescale: 48000, + AudioChannels: 2, + Data: []byte("audio-init"), + }}, + {typ: ingestframe.LLPart, event: llhls.Event{ + Kind: llhls.Part, + Presentation: "whip-1-test", + Session: 1, + Track: "video", + Generation: 1, + Timescale: 90000, + MSN: 0, + Part: 0, + Start: 0, + Duration: time.Second, + Independent: true, + Data: []byte("video-part"), + }}, + {typ: ingestframe.LLSegmentComplete, event: llhls.Event{ + Kind: llhls.SegmentComplete, + Presentation: "whip-1-test", + Session: 1, + Track: "video", + Generation: 1, + Timescale: 90000, + MSN: 0, + Start: 0, + Duration: time.Second, + Data: []byte("video-segment"), + }}, + } + + for _, test := range frames { + frame := llhlsEventToFrame(test.event) + var err error + payload, err := ingestframe.EncodeLLFrame(frame) + if err != nil { + t.Fatal(err) + } + if err := mm.observeWorkerLLHLSFrame(streamer, test.typ, payload); err != nil { + t.Fatal(err) + } + } + + config := window.VideoConfig() + if config.FrameRate != 120 { + t.Fatalf("worker video frame rate = %v, want 120", config.FrameRate) + } + if got := window.AudioConfig().Channels; got != 2 { + t.Fatalf("worker audio channels = %d, want 2", got) + } + snapshot := window.Snapshot("whip-1-test", "video") + if len(snapshot.Init) == 0 || len(snapshot.Segments) != 1 || len(snapshot.Segments[0].Parts) != 1 { + t.Fatalf("worker LL-HLS window = %+v", snapshot) + } +} + +func TestObserveWorkerLLHLSFrameRemovesWindowAtSessionEnd(t *testing.T) { + const streamer = "did:key:z6MkWorkerLLHLSEndTest" + mm := &MediaManager{} + mm.replaceLLWindow(streamer) + initEvent := llhls.Event{ + Kind: llhls.Init, + Presentation: "whip-1-test", + Session: 1, + Track: "video", + Generation: 1, + Data: []byte("video-init"), + } + frame := llhlsEventToFrame(initEvent) + var err error + payload, err := ingestframe.EncodeLLFrame(frame) + if err != nil { + t.Fatal(err) + } + if err := mm.observeWorkerLLHLSFrame(streamer, ingestframe.LLInit, payload); err != nil { + t.Fatal(err) + } + endEvent := initEvent + endEvent.Kind = llhls.SessionEnd + frame = llhlsEventToFrame(endEvent) + payload, err = ingestframe.EncodeLLFrame(frame) + if err != nil { + t.Fatal(err) + } + if err := mm.observeWorkerLLHLSFrame(streamer, ingestframe.LLSessionEnd, payload); err != nil { + t.Fatal(err) + } + if mm.GetLLWindow(streamer) != nil { + t.Fatal("LL-HLS window survived session end") + } +} + +func TestDurationFromLLTicksRejectsOverflow(t *testing.T) { + if _, err := durationFromLLTicks(^uint64(0), 1); err == nil { + t.Fatal("overflowing LL-HLS tick value was accepted") + } + if _, err := durationFromLLTicks(^uint64(0), 0); err == nil { + t.Fatal("overflowing unscaled LL-HLS tick value was accepted") + } +} diff --git a/pkg/media/llhls_pipeline.go b/pkg/media/llhls_pipeline.go new file mode 100644 index 00000000..932b52d5 --- /dev/null +++ b/pkg/media/llhls_pipeline.go @@ -0,0 +1,142 @@ +package media + +import ( + "context" + "fmt" + "time" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/log" +) + +const ( + llhlsParentDuration = 2 * time.Second + llhlsPartDuration = time.Second + llhlsPartTarget = 1100 * time.Millisecond +) + +// llhlsMuxerElements contains the shared LL-HLS CMAF muxers used by RTMP and WHIP. +// The input branches remain source-specific because WHIP must encode Opus to +// AAC while RTMP already supplies AAC. +func llhlsMuxerElements(videoChunkDuration time.Duration) []string { + if videoChunkDuration <= 0 { + videoChunkDuration = llhlsPartDuration + } + return []string{ + fmt.Sprintf("isofmp4mux name=ll_video_mux fragment-duration=%d chunk-duration=%d ! appsink name=ll_video_sink sync=false async=false", llhlsParentDuration, videoChunkDuration), + fmt.Sprintf("isofmp4mux name=ll_audio_mux manual-split=true fragment-duration=%d chunk-duration=%d ! appsink name=ll_audio_sink sync=false async=false", llhlsParentDuration, llhlsPartDuration), + } +} + +func installCMAFBranch(ctx context.Context, pipeline *gst.Pipeline, output *llhlsIngestOutput) error { + if output == nil { + return fmt.Errorf("LL-HLS CMAF sink: missing output") + } + // Both rendition playlists use the same program-date-time anchor. + programDateTimeBase := time.Now().UTC() + installTrack := func(name, track string, audioOnly bool) error { + element, err := pipeline.GetElementByName(name) + if err != nil { + return fmt.Errorf("LL-HLS CMAF sink: %w", err) + } + installCMAFSink(ctx, app.SinkFromElement(element), &cmafTrackSink{ + presentation: output.presentation, + session: output.session, + track: track, + window: output.window, + publish: output.publish, + generation: 1, + partDuration: llhlsPartDuration, + programDateTimeBase: programDateTimeBase, + partTarget: llhlsPartTarget, + audioOnly: audioOnly, + }) + return nil + } + if err := installTrack("ll_video_sink", "video", false); err != nil { + return err + } + return installTrack("ll_audio_sink", "audio", true) +} + +// startLLAudioSplitter triggers manual muxer splits from AAC buffer PTS. The +// muxer cuts between input buffers, preserving AAC frames at parent boundaries. +func startLLAudioSplitter(ctx context.Context, pipeline *gst.Pipeline) error { + mux := safeElement(pipeline, "ll_audio_mux") + if mux == nil { + return fmt.Errorf("LL-HLS audio splitter: ll_audio_mux is missing") + } + queue := safeElement(pipeline, "ll_audio_queue") + if queue == nil { + return fmt.Errorf("LL-HLS audio splitter: ll_audio_queue is missing") + } + queueSrc := queue.GetStaticPad("src") + if queueSrc == nil { + return fmt.Errorf("LL-HLS audio splitter: ll_audio_queue has no source pad") + } + splitter := &llAudioSplitter{} + probeID := queueSrc.AddProbe(gst.PadProbeTypeBuffer, func(pad *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { + buffer := info.GetBuffer() + if buffer != nil { + splitter.handleBuffer(ctx, pad, buffer) + } + return gst.PadProbeOK + }) + go func() { + <-ctx.Done() + queueSrc.RemoveProbe(probeID) + }() + return nil +} + +type llAudioSplitter struct { + initialized bool + nextBoundary time.Duration + splitIndex uint64 +} + +func (s *llAudioSplitter) handleBuffer(ctx context.Context, pad *gst.Pad, buffer *gst.Buffer) { + if buffer.GetFlags()&gst.BufferFlagDiscont != 0 { + s.initialized = false + } + pts := buffer.PresentationTimestamp().AsDuration() + if pts == nil { + return + } + if !s.initialized { + s.nextBoundary = *pts + llhlsPartDuration + s.initialized = true + return + } + if *pts < s.nextBoundary { + return + } + + chunk := (s.splitIndex+1)%2 == 1 + if err := emitLLAudioManualSplit(pad, chunk); err != nil { + log.Warn(ctx, "LL-HLS audio split event failed", "chunk", chunk, "error", err) + return + } + s.splitIndex++ + s.nextBoundary += llhlsPartDuration +} + +func emitLLAudioSplit(mux *gst.Element, boundary gst.ClockTime) error { + _, err := mux.Emit("split-at-running-time", uint64(boundary)) + return err +} + +func emitLLAudioManualSplit(queueSrc *gst.Pad, chunk bool) error { + if queueSrc == nil { + return fmt.Errorf("LL-HLS audio splitter: nil queue source pad") + } + structure := gst.NewStructure("FMP4MuxSplitNow") + if err := structure.SetValue("chunk", chunk); err != nil { + return fmt.Errorf("set audio split event: %w", err) + } + if !queueSrc.PushEvent(gst.NewCustomEvent(gst.EventTypeCustomDownstream, structure)) { + return fmt.Errorf("push audio split event") + } + return nil +} diff --git a/pkg/media/media.go b/pkg/media/media.go index b6a348ba..e52ce9fb 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -10,6 +10,7 @@ import ( "io" "sync" "sync/atomic" + "time" "github.com/google/uuid" "github.com/pion/webrtc/v4" @@ -94,7 +95,20 @@ type MediaManager struct { // session. Epochs are strictly increasing, so a newer session always wins over a // transcoder built for an older one. func (mm *MediaManager) nextIngestSession() uint64 { - return mm.ingestSessionSeq.Add(1) + // Seed the process-local sequence from wall time so a worker resumed after a + // main restart still sorts before newly-created sessions. The CAS preserves + // strict ordering when sessions start in the same nanosecond. + now := uint64(time.Now().UnixNano()) + for { + previous := mm.ingestSessionSeq.Load() + next := now + if next <= previous { + next = previous + 1 + } + if mm.ingestSessionSeq.CompareAndSwap(previous, next) { + return next + } + } } type NewSegmentNotification struct { diff --git a/pkg/media/rtmp_ingest.go b/pkg/media/rtmp_ingest.go index a191d533..54b06919 100644 --- a/pkg/media/rtmp_ingest.go +++ b/pkg/media/rtmp_ingest.go @@ -10,7 +10,6 @@ import ( "github.com/bluenviron/gortsplib/v5/pkg/format" "github.com/bluenviron/mediacommon/v2/pkg/codecs/h264" "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" "github.com/google/uuid" "stream.place/streamplace/pkg/llhls" "stream.place/streamplace/pkg/log" @@ -34,12 +33,6 @@ type RTMPSession struct { MediaSigner MediaSigner } -const ( - llhlsParentDuration = 2 * time.Second - llhlsPartDuration = time.Second - llhlsPartTarget = 1100 * time.Millisecond -) - func h264VideoConfig(track *format.H264) llhls.VideoConfig { if track == nil { return llhls.VideoConfig{} @@ -80,11 +73,10 @@ func (mm *MediaManager) RTMPIngest(ctx context.Context, rtmpURL string, ms Media llEnabled = false } if llEnabled { - // AAC and H.264 use separate renditions so each track can preserve its - // own sample boundaries. The playlists share a program-date-time grid. + // AAC and H.264 use separate muxers so each track can preserve its own + // sample boundaries. The playlists share a program-date-time grid. + pipelineSlice = append(pipelineSlice, llhlsMuxerElements(llhlsPartDuration)...) pipelineSlice = append(pipelineSlice, - fmt.Sprintf("isofmp4mux name=ll_video_mux fragment-duration=%d chunk-duration=%d ! appsink name=ll_video_sink sync=false async=false", llhlsParentDuration, llhlsPartDuration), - fmt.Sprintf("isofmp4mux name=ll_audio_mux manual-split=true fragment-duration=%d chunk-duration=%d ! appsink name=ll_audio_sink sync=false async=false", llhlsParentDuration, llhlsPartDuration), "demux.audio ! queue ! aacparse name=audioenc ! audio/mpeg,mpegversion=4,stream-format=raw ! tee name=audio_tee", // The manual-split mux waits for a future AAC buffer to carry each // split marker. The default one-second queue time limit can fill @@ -141,7 +133,11 @@ func (mm *MediaManager) RTMPIngest(ctx context.Context, rtmpURL string, ms Media window := mm.replaceLLWindow(streamerDID) defer mm.removeLLWindow(streamerDID, presentation, window) window.SetVideoConfig(h264VideoConfig(videoTrack)) - if err := installCMAFBranch(ctx, pipeline, window, presentation, session); err != nil { + if err := installCMAFBranch(ctx, pipeline, &llhlsIngestOutput{ + presentation: presentation, + session: session, + window: window, + }); err != nil { return err } } else if err = audioenc.Link(signer); err != nil { @@ -207,119 +203,3 @@ func linkElementToPad(source, destination *gst.Element, sinkPadName string) erro } return nil } - -func installCMAFBranch(ctx context.Context, pipeline *gst.Pipeline, window *llhls.Window, presentation string, session uint64) error { - // Both rendition playlists use the same program-date-time anchor. - programDateTimeBase := time.Now().UTC() - videoElement, err := pipeline.GetElementByName("ll_video_sink") - if err != nil { - return fmt.Errorf("LL-HLS CMAF sink: %w", err) - } - installCMAFSink(ctx, app.SinkFromElement(videoElement), &cmafTrackSink{ - presentation: presentation, - session: session, - track: "video", - window: window, - generation: 1, - partDuration: llhlsPartDuration, - programDateTimeBase: programDateTimeBase, - partTarget: llhlsPartTarget, - }) - audioElement, err := pipeline.GetElementByName("ll_audio_sink") - if err != nil { - return fmt.Errorf("LL-HLS CMAF audio sink: %w", err) - } - installCMAFSink(ctx, app.SinkFromElement(audioElement), &cmafTrackSink{ - presentation: presentation, - session: session, - track: "audio", - window: window, - generation: 1, - partDuration: llhlsPartDuration, - programDateTimeBase: programDateTimeBase, - partTarget: llhlsPartTarget, - audioOnly: true, - }) - return nil -} - -// startLLAudioSplitter triggers manual muxer splits from AAC buffer PTS. The -// muxer cuts between input buffers, preserving AAC frames at parent boundaries. -func startLLAudioSplitter(ctx context.Context, pipeline *gst.Pipeline) error { - mux := safeElement(pipeline, "ll_audio_mux") - if mux == nil { - return fmt.Errorf("LL-HLS audio splitter: ll_audio_mux is missing") - } - queue := safeElement(pipeline, "ll_audio_queue") - if queue == nil { - return fmt.Errorf("LL-HLS audio splitter: ll_audio_queue is missing") - } - queueSrc := queue.GetStaticPad("src") - if queueSrc == nil { - return fmt.Errorf("LL-HLS audio splitter: ll_audio_queue has no source pad") - } - splitter := &llAudioSplitter{} - probeID := queueSrc.AddProbe(gst.PadProbeTypeBuffer, func(pad *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { - buffer := info.GetBuffer() - if buffer != nil { - splitter.handleBuffer(ctx, pad, buffer) - } - return gst.PadProbeOK - }) - go func() { - <-ctx.Done() - queueSrc.RemoveProbe(probeID) - }() - return nil -} - -type llAudioSplitter struct { - initialized bool - nextBoundary time.Duration - splitIndex uint64 -} - -func (s *llAudioSplitter) handleBuffer(ctx context.Context, pad *gst.Pad, buffer *gst.Buffer) { - if buffer.GetFlags()&gst.BufferFlagDiscont != 0 { - s.initialized = false - } - pts := buffer.PresentationTimestamp().AsDuration() - if pts == nil { - return - } - if !s.initialized { - s.nextBoundary = *pts + llhlsPartDuration - s.initialized = true - return - } - if *pts < s.nextBoundary { - return - } - - chunk := (s.splitIndex+1)%2 == 1 - if err := emitLLAudioManualSplit(pad, chunk); err != nil { - log.Warn(ctx, "LL-HLS audio split event failed", "chunk", chunk, "error", err) - return - } - s.splitIndex++ - s.nextBoundary += llhlsPartDuration -} - -func emitLLAudioSplit(mux *gst.Element, boundary gst.ClockTime) error { - _, err := mux.Emit("split-at-running-time", uint64(boundary)) - return err -} - -func emitLLAudioManualSplit(queueSrc *gst.Pad, chunk bool) error { - if queueSrc == nil { - return fmt.Errorf("LL-HLS audio splitter: nil queue source pad") - } - structure := gst.NewStructure("FMP4MuxSplitNow") - if err := structure.SetValue("chunk", chunk); err != nil { - return fmt.Errorf("set audio split event: %w", err) - } - if !queueSrc.PushEvent(gst.NewCustomEvent(gst.EventTypeCustomDownstream, structure)) { - return fmt.Errorf("push audio split event") - } - return nil -} diff --git a/pkg/media/rtmp_ingest_test.go b/pkg/media/rtmp_ingest_test.go index a832bccc..b5f14281 100644 --- a/pkg/media/rtmp_ingest_test.go +++ b/pkg/media/rtmp_ingest_test.go @@ -3,6 +3,7 @@ package media import ( "context" "math" + "strings" "testing" "time" @@ -12,6 +13,18 @@ import ( "stream.place/streamplace/pkg/gstinit" ) +func TestLLHLSMuxerElementsUseSourceSpecificVideoChunkDuration(t *testing.T) { + pipeline := strings.Join(llhlsMuxerElements(500*time.Millisecond), "\n") + for _, want := range []string{ + "isofmp4mux name=ll_video_mux fragment-duration=2000000000 chunk-duration=500000000", + "isofmp4mux name=ll_audio_mux manual-split=true fragment-duration=2000000000 chunk-duration=1000000000", + } { + if !strings.Contains(pipeline, want) { + t.Fatalf("LL-HLS muxer pipeline = %q, missing %q", pipeline, want) + } + } +} + func TestH264VideoConfigUsesSPSMetadata(t *testing.T) { sps := []byte{ 0x67, 0x64, 0x00, 0x1f, 0xac, 0xd9, 0x40, 0x50, diff --git a/pkg/media/webrtc_ingest.go b/pkg/media/webrtc_ingest.go index 70b6b2bf..4cc499b6 100644 --- a/pkg/media/webrtc_ingest.go +++ b/pkg/media/webrtc_ingest.go @@ -30,7 +30,25 @@ func (mm *MediaManager) WebRTCIngest(ctx context.Context, offer *webrtc.SessionD cancel() return nil, fmt.Errorf("failed create signer element: %w", err) } - return mm.webRTCIngestPipeline(ctx, cancel, offer, peerConnection, signerElem, signer, done) + var llOutput *llhlsIngestOutput + if webRTCLLHLSAvailable(mm.cli) { + session := mm.nextIngestSession() + presentation := fmt.Sprintf("whip-%d-%s", session, uu.String()) + window := mm.replaceLLWindow(signer.Streamer()) + llOutput = &llhlsIngestOutput{ + presentation: presentation, + session: session, + window: window, + done: func() { + mm.removeLLWindow(signer.Streamer(), presentation, window) + }, + } + } + answer, err := mm.webRTCIngestPipeline(ctx, cancel, offer, peerConnection, signerElem, signer, done, llOutput) + if err != nil && llOutput != nil { + llOutput.done() + } + return answer, err } // webRTCIngestPipeline runs WebRTC ingest over a pre-built signer element: @@ -40,7 +58,7 @@ func (mm *MediaManager) WebRTCIngest(ctx context.Context, offer *webrtc.SessionD // worker passes a muxlSignSegmentElem wired to its frame socket and a nil // keyRevSigner. The cancellable ctx and signerElem are built by the caller (the // signer element's goroutines are tied to ctx). -func (mm *MediaManager) webRTCIngestPipeline(ctx context.Context, cancel context.CancelFunc, offer *webrtc.SessionDescription, peerConnection rtcrec.PeerConnection, signerElem *gst.Element, keyRevSigner MediaSigner, done chan error) (*webrtc.SessionDescription, error) { +func (mm *MediaManager) webRTCIngestPipeline(ctx context.Context, cancel context.CancelFunc, offer *webrtc.SessionDescription, peerConnection rtcrec.PeerConnection, signerElem *gst.Element, keyRevSigner MediaSigner, done chan error, llOutput *llhlsIngestOutput) (*webrtc.SessionDescription, error) { // Allow us to receive 1 audio track, and 1 video track if _, err := peerConnection.AddTransceiverFromKind(webrtc.RTPCodecTypeAudio); err != nil { return nil, fmt.Errorf("failed to add audio transceiver: %w", err) @@ -48,10 +66,24 @@ func (mm *MediaManager) webRTCIngestPipeline(ctx context.Context, cancel context return nil, fmt.Errorf("failed to add video transceiver: %w", err) } - pipelineSlice := []string{ - "multiqueue name=queue", - "appsrc format=time is-live=true do-timestamp=true name=videosrc ! capsfilter caps=application/x-rtp ! rtph264depay ! capsfilter caps=video/x-h264,stream-format=byte-stream,alignment=nal ! h264parse disable-passthrough=true config-interval=-1 ! h264timestamper ! identity ! queue.sink_0", - "appsrc format=time do-timestamp=true name=audiosrc ! capsfilter caps=application/x-rtp,media=audio,encoding-name=OPUS,payload=111 ! rtpopusdepay ! opusparse ! queue.sink_1", + pipelineSlice := []string{"multiqueue name=queue"} + videoInput := "appsrc format=time is-live=true do-timestamp=true name=videosrc ! capsfilter caps=application/x-rtp ! rtph264depay ! capsfilter caps=video/x-h264,stream-format=byte-stream,alignment=nal ! h264parse disable-passthrough=true config-interval=-1 ! h264timestamper" + audioInput := "appsrc format=time do-timestamp=true name=audiosrc ! capsfilter caps=application/x-rtp,media=audio,encoding-name=OPUS,payload=111 ! rtpopusdepay ! opusparse" + if llOutput == nil { + pipelineSlice = append(pipelineSlice, + videoInput+" ! identity ! queue.sink_0", + audioInput+" ! queue.sink_1", + ) + } else { + pipelineSlice = append(pipelineSlice, llhlsMuxerElements(llhlsPartDuration/2)...) + pipelineSlice = append(pipelineSlice, + videoInput+" ! tee name=video_tee", + "video_tee. ! queue name=video_signer_queue ! queue.sink_0", + "video_tee. ! queue ! h264parse ! video/x-h264,stream-format=avc,alignment=au ! ll_video_mux.", + audioInput+" ! tee name=audio_tee", + "audio_tee. ! queue name=audio_signer_queue ! queue.sink_1", + "audio_tee. ! queue name=ll_audio_encode_queue max-size-time=0 ! opusdec ! audioconvert ! audioresample ! fdkaacenc bitrate=128000 ! aacparse ! audio/mpeg,mpegversion=4,stream-format=raw,rate=48000,channels=2 ! queue name=ll_audio_queue max-size-time=0 ! ll_audio_mux.", + ) } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) @@ -136,6 +168,12 @@ func (mm *MediaManager) webRTCIngestPipeline(ctx context.Context, cancel context cancel() return nil, fmt.Errorf("failed to link audioSrcPad to signerElemAudioPad") } + if llOutput != nil { + if err := installCMAFBranch(ctx, pipeline, llOutput); err != nil { + cancel() + return nil, err + } + } // Setup complete! Now we boot up streaming in the background while returning the SDP offer to the user. go func() { @@ -185,6 +223,11 @@ func (mm *MediaManager) webRTCIngestPipeline(ctx context.Context, cancel context if err != nil { log.Log(ctx, "failed to set pipeline state", "error", err) cancel() + } else if llOutput != nil { + if err := startLLAudioSplitter(ctx, pipeline); err != nil { + log.Log(ctx, "failed to start LL-HLS audio splitter", "error", err) + cancel() + } } // Set the handler for ICE connection state @@ -322,6 +365,9 @@ func (mm *MediaManager) webRTCIngestPipeline(ctx context.Context, cancel context if err := videoSrcElem.SetState(gst.StateNull); err != nil { log.Log(ctx, "failed to set videoSrcElem state to null", "error", err) } + if llOutput != nil && llOutput.done != nil { + llOutput.done() + } log.Log(ctx, "webrtc ingest pipeline done") diff --git a/pkg/media/whip_worker.go b/pkg/media/whip_worker.go index 50c0465c..24584d9b 100644 --- a/pkg/media/whip_worker.go +++ b/pkg/media/whip_worker.go @@ -88,10 +88,20 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er if err != nil { return finish(fmt.Errorf("build signer element: %w", err)) } + var llOutput *llhlsIngestOutput + if cfg.LLHLS { + workerOutput := newWorkerLLHLSOutput(srv, cfg.LLHLSPresentation, cfg.LLHLSSession) + llOutput = &llhlsIngestOutput{ + presentation: cfg.LLHLSPresentation, + session: cfg.LLHLSSession, + publish: workerOutput.publish, + done: workerOutput.done, + } + } offer := &webrtc.SessionDescription{Type: webrtc.SDPTypeOffer, SDP: cfg.OfferSDP} streamDone := make(chan error, 1) - answer, err := mm.webRTCIngestPipeline(ctx, cancel, offer, pc, signerElem, nil, streamDone) + answer, err := mm.webRTCIngestPipeline(ctx, cancel, offer, pc, signerElem, nil, streamDone, llOutput) if err != nil { return finish(fmt.Errorf("webrtc ingest: %w", err)) } diff --git a/pkg/media/whip_worker_test.go b/pkg/media/whip_worker_test.go index 481e3496..cbbb4222 100644 --- a/pkg/media/whip_worker_test.go +++ b/pkg/media/whip_worker_test.go @@ -154,6 +154,9 @@ func produceWHIPMedia(t *testing.T, ctx context.Context, video, audio *webrtc.Tr func TestWHIPWorkerLoopback(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() + if gst.Find("isofmp4mux") == nil || gst.Find("fdkaacenc") == nil || gst.Find("opusdec") == nil { + t.Skip("static GStreamer build with WebRTC LL-HLS audio support is required") + } ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) @@ -174,16 +177,19 @@ func TestWHIPWorkerLoopback(t *testing.T) { sock := filepath.Join(t.TempDir(), "whip.sock") cfg := IngestWorkerConfig{ - StreamerDID: ms.Streamer(), - KeyPEM: keyPEM, - CertPEM: ms.Cert, - Manifest: manifest, - NodeCertPEM: ms.Cert, - NodeKeyPEM: keyPEM, - BroadcasterHost: "test.example.com", - SocketPath: sock, - Transport: IngestTransportWHIP, - OfferSDP: offer.SDP, + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + NodeCertPEM: ms.Cert, + NodeKeyPEM: keyPEM, + BroadcasterHost: "test.example.com", + SocketPath: sock, + Transport: IngestTransportWHIP, + OfferSDP: offer.SDP, + LLHLS: true, + LLHLSPresentation: "whip-1-test", + LLHLSSession: 1, } serveDone := make(chan error, 1) go func() { serveDone <- ServeWHIPIngestWorkerSocket(ctx, cfg) }() @@ -210,22 +216,35 @@ func TestWHIPWorkerLoopback(t *testing.T) { produceWHIPMedia(t, ctx, videoTrack, audioTrack) - // Read signed segments; require at least one valid dual-codec one. + // Read LL-HLS frames and signed segments; require both media tracks and at + // least one valid dual-codec signed segment. _ = conn.SetReadDeadline(time.Now().Add(45 * time.Second)) - var segs int - for segs == 0 { + var segs, llInits, llParts int + for segs == 0 || llInits < 2 || llParts == 0 { typ, payload, ferr := fr.ReadFrame() require.NoError(t, ferr, "reading worker frames") - if typ != ingestframe.Segment { + switch typ { + case ingestframe.LLInit: + frame, derr := ingestframe.DecodeLLFrame(payload) + require.NoError(t, derr) + require.Equal(t, "whip-1-test", frame.Presentation) + llInits++ + case ingestframe.LLPart: + frame, derr := ingestframe.DecodeLLFrame(payload) + require.NoError(t, derr) + require.Equal(t, "whip-1-test", frame.Presentation) + llParts++ + case ingestframe.Segment: + out, verr := muxl.RunMuxlVerify(ctx, bytes.NewReader(payload)) + require.NoError(t, verr) + require.NotContains(t, out, `"validation_state":"Invalid"`, "segment must validate") + segs++ + default: continue } - out, verr := muxl.RunMuxlVerify(ctx, bytes.NewReader(payload)) - require.NoError(t, verr) - require.NotContains(t, out, `"validation_state":"Invalid"`, "segment must validate") - segs++ } _ = conn.SetReadDeadline(time.Time{}) - t.Logf("whip worker produced %d signed segment(s) from real RTP media", segs) + t.Logf("whip worker produced %d init frame(s), %d part frame(s), and %d signed segment(s) from real RTP media", llInits, llParts, segs) // Close the connection before tearing down: in production main's connection // breaks on shutdown, which detaches the worker's frame server so it drains