diff --git a/pkg/api/playback_llhls.go b/pkg/api/playback_llhls.go index 0a308dd8..36667e2c 100644 --- a/pkg/api/playback_llhls.go +++ b/pkg/api/playback_llhls.go @@ -46,23 +46,37 @@ func (a *StreamplaceAPI) HandleLLHLSMaster(ctx context.Context) httprouter.Handl http.Redirect(w, r, "/xrpc/place.stream.playback.getLivePlaylist?streamer="+url.QueryEscape(p.ByName("user")), http.StatusTemporaryRedirect) return } + if len(window.Snapshot(presentation, "video").Init) == 0 || len(window.Snapshot(presentation, "audio").Init) == 0 { + w.Header().Set("Cache-Control", "no-store") + http.Redirect(w, r, "/xrpc/place.stream.playback.getLivePlaylist?streamer="+url.QueryEscape(p.ByName("user")), http.StatusTemporaryRedirect) + return + } base := fmt.Sprintf("/api/playback/%s/llhls/%s", url.PathEscape(did), url.PathEscape(presentation)) - body := renderLLHLSMaster(base, window.VideoConfig()) + body := renderLLHLSMaster(base, window.VideoConfig(), window.AudioConfig()) writeLLHLSPlaylist(w, body) } } -func renderLLHLSMaster(base string, videoConfig llhls.VideoConfig) string { +func renderLLHLSMaster(base string, videoConfig llhls.VideoConfig, audioConfig llhls.AudioConfig) string { codec := videoConfig.Codec if codec == "" { codec = "avc1.64001f" } - streamInf := fmt.Sprintf("#EXT-X-STREAM-INF:BANDWIDTH=2500000,CODECS=%q", codec+",mp4a.40.2") + channels := audioConfig.Channels + if channels <= 0 { + channels = 2 + } + var b strings.Builder + b.WriteString("#EXTM3U\n#EXT-X-VERSION:10\n") + fmt.Fprintf(&b, "#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID=%q,NAME=%q,DEFAULT=YES,AUTOSELECT=YES,CHANNELS=%q,CODECS=%q,URI=%q\n", + "audio", "default", strconv.Itoa(channels), "mp4a.40.2", base+"/audio/index.m3u8") + streamInf := fmt.Sprintf("#EXT-X-STREAM-INF:BANDWIDTH=6500000,CODECS=%q", codec+",mp4a.40.2") if videoConfig.Width > 0 && videoConfig.Height > 0 { streamInf += fmt.Sprintf(",RESOLUTION=%dx%d", videoConfig.Width, videoConfig.Height) } - streamInf += ",CLOSED-CAPTIONS=NONE" - return "#EXTM3U\n#EXT-X-VERSION:10\n" + streamInf + "\n" + base + "/video/index.m3u8\n" + streamInf += ",AUDIO=\"audio\",CLOSED-CAPTIONS=NONE" + b.WriteString(streamInf + "\n" + base + "/video/index.m3u8\n") + return b.String() } func (a *StreamplaceAPI) HandleLLHLS(ctx context.Context) httprouter.Handle { @@ -136,10 +150,19 @@ func (a *StreamplaceAPI) HandleLLHLSPlaylist(ctx context.Context) httprouter.Han } } base := fmt.Sprintf("/api/playback/%s/llhls/%s/%s", url.PathEscape(did), url.PathEscape(presentation), track) - body := window.Playlist(presentation, track, + renditionBase := fmt.Sprintf("/api/playback/%s/llhls/%s", url.PathEscape(did), url.PathEscape(presentation)) + playlist := window.Playlist + if track == "audio" { + // Safari's native LL-HLS path loses audio when AAC parts are exposed + // as a second low-latency rendition. Keep audio on complete parents; + // video remains low latency. + playlist = window.PlaylistSegmentsOnly + } + body := playlist(presentation, track, func(msn uint64, part uint32) string { return fmt.Sprintf("%s/%d/%d.m4s", base, msn, part) }, func(msn uint64) string { return fmt.Sprintf("%s/%d.m4s", base, msn) }, - base+"/init.mp4") + base+"/init.mp4", + func(otherTrack string) string { return fmt.Sprintf("%s/%s/index.m3u8", renditionBase, otherTrack) }) if body == "" { apierrors.WriteHTTPNotFound(w, "track not found", nil) return @@ -193,7 +216,7 @@ func (a *StreamplaceAPI) HandleLLHLSInit(ctx context.Context) httprouter.Handle } else { log.Debug(r.Context(), "LL-HLS init response", "presentation", presentation, "track", track, "bytes", len(data)) } - serveLLHLSBytes(w, r, data, "init.mp4", false) + serveLLHLSBytes(w, r, data, "init.mp4", false, llhlsTrackMIMEType(track)) } } @@ -222,7 +245,7 @@ func (a *StreamplaceAPI) HandleLLHLSPart(ctx context.Context) httprouter.Handle } else { log.Debug(r.Context(), "LL-HLS part response", "presentation", presentation, "track", track, "msn", msn, "part", partIndex, "bytes", len(data)) } - serveLLHLSBytes(w, r, data, "part.m4s", true) + serveLLHLSBytes(w, r, data, "part.m4s", true, llhlsTrackMIMEType(track)) } } @@ -245,7 +268,7 @@ func (a *StreamplaceAPI) HandleLLHLSSegment(ctx context.Context) httprouter.Hand } else { log.Debug(r.Context(), "LL-HLS segment response", "presentation", presentation, "track", track, "msn", msn, "bytes", len(data)) } - serveLLHLSBytes(w, r, data, "segment.m4s", true) + serveLLHLSBytes(w, r, data, "segment.m4s", true, llhlsTrackMIMEType(track)) } } @@ -256,13 +279,20 @@ func writeLLHLSPlaylist(w http.ResponseWriter, body string) { _, _ = w.Write([]byte(body)) } -func serveLLHLSBytes(w http.ResponseWriter, r *http.Request, data []byte, name string, immutable bool) { +func llhlsTrackMIMEType(track string) string { + if track == "audio" { + return "audio/mp4" + } + return "video/mp4" +} + +func serveLLHLSBytes(w http.ResponseWriter, r *http.Request, data []byte, name string, immutable bool, contentType string) { if len(data) == 0 { w.Header().Set("Cache-Control", "no-store") apierrors.WriteHTTPNotFound(w, "media not found", nil) return } - w.Header().Set("Content-Type", "video/mp4") + w.Header().Set("Content-Type", contentType) if immutable { w.Header().Set("Cache-Control", "public, max-age=31536000, immutable") } else { diff --git a/pkg/api/playback_llhls_test.go b/pkg/api/playback_llhls_test.go index 62c4256e..56825757 100644 --- a/pkg/api/playback_llhls_test.go +++ b/pkg/api/playback_llhls_test.go @@ -55,32 +55,37 @@ func TestLLHLSMasterAdvertisesVideoMetadata(t *testing.T) { Codec: "avc1.64002a", Width: 1280, Height: 720, - }) + }, llhls.AudioConfig{Channels: 2}) for _, want := range []string{ "#EXT-X-VERSION:10", - `#EXT-X-STREAM-INF:BANDWIDTH=2500000,CODECS="avc1.64002a,mp4a.40.2",RESOLUTION=1280x720,CLOSED-CAPTIONS=NONE`, + `#EXT-X-STREAM-INF:BANDWIDTH=6500000,CODECS="avc1.64002a,mp4a.40.2",RESOLUTION=1280x720,AUDIO="audio",CLOSED-CAPTIONS=NONE`, + `#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID="audio",NAME="default",DEFAULT=YES,AUTOSELECT=YES,CHANNELS="2",CODECS="mp4a.40.2",URI="/api/playback/did:plc:test/llhls/rtmp-1/audio/index.m3u8"`, `/api/playback/did:plc:test/llhls/rtmp-1/video/index.m3u8`, } { if !strings.Contains(master, want) { t.Errorf("master missing %q:\n%s", want, master) } } - if strings.Contains(master, "#EXT-X-MEDIA:TYPE=AUDIO") || strings.Contains(master, `AUDIO="audio"`) { - t.Fatalf("muxed master should not advertise a separate audio rendition:\n%s", master) - } if strings.Contains(master, "#EXT-X-INDEPENDENT-SEGMENTS") { t.Fatalf("RTMP master cannot make a presentation-wide independence guarantee:\n%s", master) } } func TestLLHLSMasterOmitsIndependentSegments(t *testing.T) { - master := renderLLHLSMaster("/api/playback/test", llhls.VideoConfig{}) + master := renderLLHLSMaster("/api/playback/test", llhls.VideoConfig{}, llhls.AudioConfig{Channels: 2}) if strings.Contains(master, "#EXT-X-INDEPENDENT-SEGMENTS") { t.Fatalf("master advertised independent segments without metadata:\n%s", master) } } +func TestLLHLSMasterAdvertisesAudioChannels(t *testing.T) { + master := renderLLHLSMaster("/api/playback/test", llhls.VideoConfig{}, llhls.AudioConfig{Channels: 1}) + if !strings.Contains(master, `CHANNELS="1"`) { + t.Fatalf("master omitted mono audio metadata:\n%s", master) + } +} + func TestLLHLSMasterRedirectsWhilePresentationIsInitializing(t *testing.T) { const user = "did:key:z6MkPreInitTest" manager := &media.MediaManager{} @@ -104,6 +109,76 @@ func TestLLHLSMasterRedirectsWhilePresentationIsInitializing(t *testing.T) { } } +func TestLLHLSMasterWaitsForBothRenditions(t *testing.T) { + const user = "did:key:z6MkBothTracksTest" + window := llhls.NewWindow() + if err := window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "p", Track: "video", Generation: 1, Data: []byte("video-init")}); err != nil { + t.Fatal(err) + } + manager := &media.MediaManager{} + setLLWindowsForTest(manager, map[string]*llhls.Window{user: window}) + api := &StreamplaceAPI{MediaManager: manager, Aliases: map[string]string{}} + handler := api.HandleLLHLSMaster(context.Background()) + + recorder := httptest.NewRecorder() + handler(recorder, httptest.NewRequest(http.MethodGet, "/master.m3u8", nil), httprouter.Params{{Key: "user", Value: user}}) + if recorder.Code != http.StatusTemporaryRedirect { + t.Fatalf("master with video-only init status = %d, want %d", recorder.Code, http.StatusTemporaryRedirect) + } + if got := recorder.Header().Get("Cache-Control"); got != "no-store" { + t.Fatalf("master with video-only init cache policy = %q, want no-store", got) + } + + if err := window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "p", Track: "audio", Generation: 1, Data: []byte("audio-init")}); err != nil { + t.Fatal(err) + } + recorder = httptest.NewRecorder() + handler(recorder, httptest.NewRequest(http.MethodGet, "/master.m3u8", nil), httprouter.Params{{Key: "user", Value: user}}) + if recorder.Code != http.StatusOK || !strings.Contains(recorder.Body.String(), `AUDIO="audio"`) { + t.Fatalf("master after both init segments = status %d body %q", recorder.Code, recorder.Body.String()) + } +} + +func TestLLHLSAudioPlaylistUsesCompleteSegments(t *testing.T) { + const user = "did:key:z6MkAudioSegmentsOnlyTest" + window := llhls.NewWindow() + if err := window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "p", Track: "audio", Generation: 1, Data: []byte("audio-init")}); err != nil { + t.Fatal(err) + } + if err := window.Observe(llhls.Event{Kind: llhls.Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Part: 0, Duration: time.Second, Data: []byte("part-0")}); err != nil { + t.Fatal(err) + } + if err := window.Observe(llhls.Event{Kind: llhls.Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Part: 1, Duration: time.Second, Data: []byte("part-1")}); err != nil { + t.Fatal(err) + } + if err := window.Observe(llhls.Event{Kind: llhls.SegmentComplete, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Duration: 2 * time.Second, Data: []byte("segment")}); err != nil { + t.Fatal(err) + } + + manager := &media.MediaManager{} + setLLWindowsForTest(manager, map[string]*llhls.Window{user: window}) + api := &StreamplaceAPI{MediaManager: manager, Aliases: map[string]string{}} + params := httprouter.Params{ + {Key: "user", Value: user}, + {Key: "presentation", Value: "p"}, + {Key: "track", Value: "audio"}, + } + recorder := httptest.NewRecorder() + api.HandleLLHLSPlaylist(context.Background())(recorder, httptest.NewRequest(http.MethodGet, "/audio/index.m3u8", nil), params) + if recorder.Code != http.StatusOK { + t.Fatalf("audio playlist status = %d, want %d", recorder.Code, http.StatusOK) + } + body := recorder.Body.String() + for _, forbidden := range []string{"#EXT-X-PART-INF:", "#EXT-X-SERVER-CONTROL:", "#EXT-X-PART:", "#EXT-X-PRELOAD-HINT:", "#EXT-X-RENDITION-REPORT:"} { + if strings.Contains(body, forbidden) { + t.Errorf("audio playlist contains %q:\n%s", forbidden, body) + } + } + if !strings.Contains(body, "#EXTINF:2.000000,") || !strings.Contains(body, "/audio/1.m4s") { + t.Fatalf("audio playlist omitted complete segment:\n%s", body) + } +} + func TestLLHLSPartHandlerKeepsExactURIIdentity(t *testing.T) { const ( user = "did:key:z6MkTest" @@ -152,6 +227,37 @@ func TestLLHLSPartHandlerKeepsExactURIIdentity(t *testing.T) { } } +func TestLLHLSPartHandlerUsesAudioMIMEType(t *testing.T) { + const user = "did:key:z6MkAudioMIMETypeTest" + window := llhls.NewWindow() + if err := window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "p", Track: "audio", Generation: 1}); err != nil { + t.Fatal(err) + } + if err := window.Observe(llhls.Event{Kind: llhls.Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Part: 0, Data: []byte("audio-part")}); err != nil { + t.Fatal(err) + } + + manager := &media.MediaManager{} + setLLWindowsForTest(manager, map[string]*llhls.Window{user: window}) + api := &StreamplaceAPI{MediaManager: manager, Aliases: map[string]string{}} + params := httprouter.Params{ + {Key: "user", Value: user}, + {Key: "presentation", Value: "p"}, + {Key: "track", Value: "audio"}, + {Key: "msn", Value: "1"}, + {Key: "part.m4s", Value: "0.m4s"}, + } + + recorder := httptest.NewRecorder() + api.HandleLLHLSPart(context.Background())(recorder, httptest.NewRequest(http.MethodGet, "/part.m4s", nil), params) + if recorder.Code != http.StatusOK { + t.Fatalf("audio part status = %d, want %d", recorder.Code, http.StatusOK) + } + if got := recorder.Header().Get("Content-Type"); got != "audio/mp4" { + t.Fatalf("audio part content type = %q, want audio/mp4", got) + } +} + func TestLLHLSPartHandlerWaitsForExactPartThenServesIt(t *testing.T) { const user = "did:key:z6MkWaitTest" window := llhls.NewWindow() @@ -202,7 +308,7 @@ func TestMissingLLHLSMediaIsNotCached(t *testing.T) { recorder := httptest.NewRecorder() request := httptest.NewRequest(http.MethodGet, "/part.m4s", nil) - serveLLHLSBytes(recorder, request, nil, "part.m4s", true) + serveLLHLSBytes(recorder, request, nil, "part.m4s", true, "video/mp4") if recorder.Code != http.StatusNotFound { t.Fatalf("missing media status = %d, want %d", recorder.Code, http.StatusNotFound) diff --git a/pkg/llhls/window.go b/pkg/llhls/window.go index cabcb27a..24903013 100644 --- a/pkg/llhls/window.go +++ b/pkg/llhls/window.go @@ -18,6 +18,7 @@ var ( ErrPartOrder = errors.New("llhls: invalid part order") ErrGeneration = errors.New("llhls: stale configuration generation") ErrPartUnavailable = errors.New("llhls: part unavailable") + ErrWindowCapacity = errors.New("llhls: window capacity exceeded") ) const ( @@ -27,6 +28,10 @@ const ( defaultTargetDuration = 6 * time.Second defaultPartTarget = 1100 * time.Millisecond partHoldBackMargin = time.Millisecond + // Eviction keeps at least this many completed segments per track so every + // rendition retains enough history for hold-back playback, even when a + // large sibling track drives the shared byte budget over its limit. + minRetainedSegments = 12 ) type EventKind uint8 @@ -66,6 +71,8 @@ type SegmentSnapshot struct { Start, Duration time.Duration Parts []PartSnapshot Complete bool + Independent bool + ProgramDateTime time.Time Data []byte } @@ -85,6 +92,10 @@ type VideoConfig struct { Height int } +type AudioConfig struct { + Channels int +} + // PartIdentity is the stable logical identity used by a part URI. Parent // timing metadata may be corrected when a parent closes, but this identity is // never reassigned or aliased to another parent's bytes. @@ -102,12 +113,14 @@ type Window struct { bytes int changed chan struct{} videoConfig VideoConfig + audioConfig AudioConfig programDateTime time.Time programDateTimeStart time.Duration targetDuration int64 partTarget time.Duration configuredTarget int64 configuredPartTarget time.Duration + completionHold time.Duration } type track struct { @@ -115,6 +128,7 @@ type track struct { init []byte segments []*segment ended bool + bytes int } type segment struct { @@ -122,6 +136,9 @@ type segment struct { start, duration time.Duration parts []*part complete bool + closing bool + independent bool + programDateTime time.Time data []byte } @@ -137,6 +154,17 @@ type Option func(*Window) func WithMaxSegments(n int) Option { return func(w *Window) { w.maxSegments = n } } func WithMaxBytes(n int) Option { return func(w *Window) { w.maxBytes = n } } +// WithSegmentCompletionDelay holds a parent segment's completion for the +// given duration after its SegmentComplete event before it becomes visible +// in snapshots and playlists. LL-HLS players must see the final part listed +// in an open segment to fetch it; when the final part and the completion +// land in the same playlist update, blocking reloads for that part roll over +// to the completed parent and players fall back to re-fetching the whole +// segment. Zero completes synchronously. +func WithSegmentCompletionDelay(d time.Duration) Option { + return func(w *Window) { w.completionHold = d } +} + // WithPlaylistDurations sets fixed upper bounds for parent and part durations. // The parent bound is rounded to the nearest whole second for TARGETDURATION. // Both values remain fixed for each presentation observed by the Window. @@ -187,6 +215,7 @@ func (w *Window) Observe(ev Event) error { } w.tracks = make(map[string]*track) w.presentation = ev.Presentation + w.audioConfig = AudioConfig{} w.programDateTime = time.Time{} w.programDateTimeStart = 0 w.targetDuration = w.configuredTarget @@ -221,11 +250,18 @@ func (w *Window) Observe(ev Event) error { if ev.Part != 0 { return ErrPartOrder } + w.completeClosingSegments(t, ev.MSN) s = w.findOrCreateSegment(t, ev) } if ev.Part != uint32(len(s.parts)) { return ErrPartOrder } + if ev.Part == 0 { + s.independent = ev.Independent + if !ev.ProgramDateTime.IsZero() { + s.programDateTime = ev.ProgramDateTime + } + } if w.programDateTime.IsZero() && !ev.ProgramDateTime.IsZero() { w.programDateTime = ev.ProgramDateTime w.programDateTimeStart = ev.Start @@ -233,19 +269,29 @@ func (w *Window) Observe(ev Event) error { p := &part{identity: PartIdentity{MSN: ev.MSN, Index: ev.Part}, start: ev.Start, duration: ev.Duration, independent: ev.Independent, data: append([]byte(nil), ev.Data...)} s.parts = append(s.parts, p) w.bytes += len(p.data) + t.bytes += len(p.data) case SegmentComplete: s := w.findSegment(t, ev.MSN) if s == nil { return fmt.Errorf("llhls: segment %d has no parts", ev.MSN) } - if s.complete { + if s.complete || s.closing { return nil } - s.start, s.duration, s.complete = ev.Start, ev.Duration, true + s.start, s.duration = ev.Start, ev.Duration s.data = append([]byte(nil), ev.Data...) w.bytes += len(s.data) + t.bytes += len(s.data) + if w.completionHold <= 0 { + s.complete = true + } else { + s.closing = true + w.scheduleCompletion(t, ev.Track, s) + } case Discontinuity: + w.removeTrackBytes(t) t.segments = nil + t.bytes = 0 t.ended = false case SessionEnd: t.ended = true @@ -254,6 +300,9 @@ func (w *Window) Observe(ev Event) error { } w.evict() w.notify() + if w.maxBytes > 0 && w.hasOverCapacityOpenSegment() { + return ErrWindowCapacity + } return nil } @@ -276,27 +325,89 @@ func (w *Window) findOrCreateSegment(t *track, ev Event) *segment { return s } +// scheduleCompletion flips a closing segment to complete after the hold. The +// timer re-validates under the lock: a presentation reset or discontinuity +// may have dropped the segment (or the whole track) in the meantime. +func (w *Window) scheduleCompletion(t *track, trackID string, s *segment) { + time.AfterFunc(w.completionHold, func() { + w.mu.Lock() + defer w.mu.Unlock() + if w.tracks[trackID] != t { + return + } + if w.findSegment(t, s.msn) != s || s.complete { + return + } + s.closing = false + s.complete = true + w.evict() + w.notify() + }) +} + +func (w *Window) completeClosingSegments(t *track, nextMSN uint64) { + for _, s := range t.segments { + if s.closing && s.msn < nextMSN { + s.closing = false + s.complete = true + } + } +} + func (w *Window) evict() { for _, t := range w.tracks { - for (w.maxSegments > 0 && len(t.segments) > w.maxSegments) || (w.maxBytes > 0 && w.bytes > w.maxBytes) { - if len(t.segments) == 0 { + for len(t.segments) > 0 { + // Incomplete segments are still receiving part events; evicting + // one strands those events and fails the ingest stream. Small + // tracks would otherwise be emptied wholesale while a large + // track keeps the byte budget exceeded. + if !t.segments[0].complete { + break + } + overCount := w.maxSegments > 0 && len(t.segments) > w.maxSegments + overBytes := w.maxBytes > 0 && w.bytes > w.maxBytes + if !overCount && !overBytes { break } - w.removeSegmentBytes(t.segments[0]) + // Retained history for hold-back playback is only shed when this + // track alone exceeds the byte budget. + if !overCount && len(t.segments) <= minRetainedSegments && t.bytes <= w.maxBytes { + break + } + w.removeSegmentBytes(t, t.segments[0]) t.segments = t.segments[1:] } } } func (w *Window) removeTrackBytes(t *track) { for _, s := range t.segments { - w.removeSegmentBytes(s) + w.removeSegmentBytes(t, s) } + t.bytes = 0 +} +func (w *Window) removeSegmentBytes(t *track, s *segment) { + bytes := segmentBytes(s) + w.bytes -= bytes + t.bytes -= bytes } -func (w *Window) removeSegmentBytes(s *segment) { + +func segmentBytes(s *segment) int { + bytes := len(s.data) for _, p := range s.parts { - w.bytes -= len(p.data) + bytes += len(p.data) } - w.bytes -= len(s.data) + return bytes +} + +func (w *Window) hasOverCapacityOpenSegment() bool { + for _, t := range w.tracks { + for _, s := range t.segments { + if !s.complete && !s.closing && segmentBytes(s) > w.maxBytes { + return true + } + } + } + return false } func (w *Window) notify() { close(w.changed); w.changed = make(chan struct{}) } @@ -321,6 +432,18 @@ func (w *Window) VideoConfig() VideoConfig { return w.videoConfig } +func (w *Window) SetAudioConfig(config AudioConfig) { + w.mu.Lock() + defer w.mu.Unlock() + w.audioConfig = config +} + +func (w *Window) AudioConfig() AudioConfig { + w.mu.Lock() + defer w.mu.Unlock() + return w.audioConfig +} + func (w *Window) Snapshot(presentation, trackID string) Snapshot { w.mu.Lock() defer w.mu.Unlock() @@ -341,7 +464,7 @@ func (w *Window) Snapshot(presentation, trackID string) Snapshot { ProgramDateTimeStart: w.programDateTimeStart, } for _, seg := range t.segments { - ss := SegmentSnapshot{MSN: seg.msn, Start: seg.start, Duration: seg.duration, Complete: seg.complete, Data: append([]byte(nil), seg.data...)} + ss := SegmentSnapshot{MSN: seg.msn, Start: seg.start, Duration: seg.duration, Complete: seg.complete, Independent: seg.independent, ProgramDateTime: seg.programDateTime, Data: append([]byte(nil), seg.data...)} for _, p := range seg.parts { ss.Parts = append(ss.Parts, PartSnapshot{Identity: p.identity, Index: p.identity.Index, Start: p.start, Duration: p.duration, Independent: p.independent, Data: append([]byte(nil), p.data...)}) } @@ -511,7 +634,20 @@ func (w *Window) partStateLocked(presentation, trackID string, msn uint64, partI // Playlist renders the media playlist for one track. URIs are supplied by the // caller so routing and presentation identifiers remain outside this package. -func (w *Window) Playlist(presentation, trackID string, partURI func(uint64, uint32) string, segmentURI func(uint64) string, initURI string) string { +// When renditionURI is non-nil, an EXT-X-RENDITION-REPORT is emitted for +// every other track that has published media (required for LL-HLS). +func (w *Window) Playlist(presentation, trackID string, partURI func(uint64, uint32) string, segmentURI func(uint64) string, initURI string, renditionURI func(string) string) string { + return w.playlist(presentation, trackID, partURI, segmentURI, initURI, renditionURI, true) +} + +// PlaylistSegmentsOnly renders a conventional media playlist using only +// completed parent segments. It is useful for renditions whose codec timing +// is not safe to expose through the LL-HLS part path. +func (w *Window) PlaylistSegmentsOnly(presentation, trackID string, partURI func(uint64, uint32) string, segmentURI func(uint64) string, initURI string, renditionURI func(string) string) string { + return w.playlist(presentation, trackID, partURI, segmentURI, initURI, renditionURI, false) +} + +func (w *Window) playlist(presentation, trackID string, partURI func(uint64, uint32) string, segmentURI func(uint64) string, initURI string, renditionURI func(string) string, includeParts bool) string { s := w.Snapshot(presentation, trackID) if s.Track == "" { return "" @@ -519,7 +655,25 @@ func (w *Window) Playlist(presentation, trackID string, partURI func(uint64, uin targetSeconds, partTarget := w.playlistDurations() partHoldBack := 3*partTarget + partHoldBackMargin var b strings.Builder - fmt.Fprintf(&b, "#EXTM3U\n#EXT-X-VERSION:10\n#EXT-X-TARGETDURATION:%d\n#EXT-X-PART-INF:PART-TARGET=%.6f\n#EXT-X-SERVER-CONTROL:CAN-BLOCK-RELOAD=YES,PART-HOLD-BACK=%.6f,HOLD-BACK=%.6f\n", targetSeconds, partTarget.Seconds(), partHoldBack.Seconds(), 3*float64(targetSeconds)) + fmt.Fprintf(&b, "#EXTM3U\n#EXT-X-VERSION:10\n#EXT-X-TARGETDURATION:%d\n", targetSeconds) + if includeParts { + fmt.Fprintf(&b, "#EXT-X-PART-INF:PART-TARGET=%.6f\n", partTarget.Seconds()) + } + if trackID == "video" && allSegmentsIndependent(s.Segments) { + b.WriteString("#EXT-X-INDEPENDENT-SEGMENTS\n") + } + if includeParts { + fmt.Fprintf(&b, "#EXT-X-SERVER-CONTROL:CAN-BLOCK-RELOAD=YES,PART-HOLD-BACK=%.6f,HOLD-BACK=%.6f\n", partHoldBack.Seconds(), 3*float64(targetSeconds)) + } + if includeParts && renditionURI != nil { + for _, rep := range w.renditionReports(presentation, trackID) { + fmt.Fprintf(&b, "#EXT-X-RENDITION-REPORT:URI=%q,LAST-MSN=%d", renditionURI(rep.trackID), rep.lastMSN) + if rep.lastPart >= 0 { + fmt.Fprintf(&b, ",LAST-PART=%d", rep.lastPart) + } + b.WriteByte('\n') + } + } fmt.Fprintf(&b, "#EXT-X-MEDIA-SEQUENCE:%d\n#EXT-X-MAP:URI=%q\n", firstMSN(s), initURI) openMSN := uint64(0) if len(s.Segments) > 0 { @@ -528,7 +682,7 @@ func (w *Window) Playlist(presentation, trackID string, partURI func(uint64, uin } } for _, seg := range s.Segments { - if !seg.Complete && seg.MSN == openMSN { + if includeParts && !seg.Complete && seg.MSN == openMSN { for _, p := range seg.Parts { fmt.Fprintf(&b, "#EXT-X-PART:DURATION=%.6f,URI=%q", p.Duration.Seconds(), partURI(p.Identity.MSN, p.Identity.Index)) if p.Independent { @@ -538,14 +692,17 @@ func (w *Window) Playlist(presentation, trackID string, partURI func(uint64, uin } } if seg.Complete { - if !s.ProgramDateTime.IsZero() { - programDateTime := s.ProgramDateTime.Add(seg.Start - s.ProgramDateTimeStart) + programDateTime := seg.ProgramDateTime + if programDateTime.IsZero() && !s.ProgramDateTime.IsZero() { + programDateTime = s.ProgramDateTime.Add(seg.Start - s.ProgramDateTimeStart) + } + if !programDateTime.IsZero() { fmt.Fprintf(&b, "#EXT-X-PROGRAM-DATE-TIME:%s\n", programDateTime.Format(time.RFC3339Nano)) } fmt.Fprintf(&b, "#EXTINF:%.6f,\n%s\n", seg.Duration.Seconds(), segmentURI(seg.MSN)) } } - if !s.Ended && len(s.Segments) > 0 { + if includeParts && !s.Ended && len(s.Segments) > 0 { last := s.Segments[len(s.Segments)-1] nextMSN, nextPart := last.MSN, uint32(len(last.Parts)) if last.Complete { @@ -560,6 +717,18 @@ func (w *Window) Playlist(presentation, trackID string, partURI func(uint64, uin return b.String() } +func allSegmentsIndependent(segments []SegmentSnapshot) bool { + if len(segments) == 0 { + return false + } + for _, segment := range segments { + if !segment.Independent { + return false + } + } + return true +} + func (w *Window) playlistDurations() (targetSeconds int64, partTarget time.Duration) { w.mu.Lock() targetSeconds, partTarget = w.targetDuration, w.partTarget @@ -567,6 +736,45 @@ func (w *Window) playlistDurations() (targetSeconds int64, partTarget time.Durat return targetSeconds, partTarget } +type renditionReport struct { + trackID string + lastMSN uint64 + lastPart int32 +} + +// renditionReports describes the latest published state of every track other +// than exclude. A track with an open segment reports that segment's MSN and +// its last published part; a fully completed tail reports its last MSN +// without a part (completed parents carry no listed partial segments). +func (w *Window) renditionReports(presentation, exclude string) []renditionReport { + w.mu.Lock() + defer w.mu.Unlock() + if presentation != w.presentation { + return nil + } + var reports []renditionReport + for trackID, t := range w.tracks { + if trackID == exclude || len(t.segments) == 0 { + continue + } + last := t.segments[len(t.segments)-1] + if !last.complete { + if len(last.parts) == 0 { + if len(t.segments) < 2 { + continue + } + last = t.segments[len(t.segments)-2] + } else { + reports = append(reports, renditionReport{trackID: trackID, lastMSN: last.msn, lastPart: int32(len(last.parts) - 1)}) + continue + } + } + reports = append(reports, renditionReport{trackID: trackID, lastMSN: last.msn, lastPart: -1}) + } + sort.Slice(reports, func(i, j int) bool { return reports[i].trackID < reports[j].trackID }) + return reports +} + func roundedDurationSeconds(d time.Duration) int64 { if d <= 0 { return 1 diff --git a/pkg/llhls/window_test.go b/pkg/llhls/window_test.go index 38690555..6e7a6c0c 100644 --- a/pkg/llhls/window_test.go +++ b/pkg/llhls/window_test.go @@ -44,6 +44,123 @@ func TestWindowPublishesPartsAndOnlyCompletesParentsWhenAllPartsArrive(t *testin } } +func completionHoldURI(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) } +func completionHoldSegment(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) } + +func waitForSegmentComplete(t *testing.T, w *Window, presentation, track string, msn uint64) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + for _, seg := range w.Snapshot(presentation, track).Segments { + if seg.MSN == msn && seg.Complete { + return + } + } + time.Sleep(5 * time.Millisecond) + } + t.Fatalf("segment %d did not complete after hold", msn) +} + +func TestWindowCompletionHoldKeepsFinalPartListed(t *testing.T) { + w := NewWindow(WithSegmentCompletionDelay(50 * time.Millisecond)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Part: 0, Start: 0, Duration: time.Second, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Part: 1, Start: time.Second, Duration: time.Second, Data: []byte("b")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Start: 0, Duration: 2 * time.Second, Data: []byte("ab")}) + + // During the hold the parent still reads as open: the final part stays + // listed and fetchable so a blocking reload for it resolves with the part + // instead of rolling over to the completed segment. + s := w.Snapshot("p", "v") + if len(s.Segments) != 1 || s.Segments[0].Complete { + t.Fatalf("segment completed during hold: %+v", s.Segments) + } + playlist := w.Playlist("p", "v", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "3/1.m4s") { + t.Fatalf("final part not listed during hold:\n%s", playlist) + } + if strings.Contains(playlist, "#EXTINF") { + t.Fatalf("segment listed as complete during hold:\n%s", playlist) + } + if err := w.Wait(context.Background(), "p", "v", 3, 1); err != nil { + t.Fatalf("blocking reload for held part: %v", err) + } + if err := w.WaitForPart(context.Background(), "p", "v", 3, 1); err != nil { + t.Fatalf("held part not fetchable: %v", err) + } + if got := w.Data("p", "v", 3, 1); !bytes.Equal(got, []byte("b")) { + t.Fatalf("held part data = %q", got) + } + + waitForSegmentComplete(t, w, "p", "v", 3) + playlist = w.Playlist("p", "v", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXTINF") || strings.Contains(playlist, "3/1.m4s") { + t.Fatalf("completed segment not published after hold:\n%s", playlist) + } + if got := w.SegmentData("p", "v", 3); !bytes.Equal(got, []byte("ab")) { + t.Fatalf("segment data = %q", got) + } + // The part stays fetchable after completion; only eviction removes it. + if err := w.WaitForPart(context.Background(), "p", "v", 3, 1); err != nil { + t.Fatalf("part after completion = %v, want still fetchable", err) + } +} + +func TestWindowCompletionHoldClosesWhenNextParentStarts(t *testing.T) { + w := NewWindow(WithSegmentCompletionDelay(50 * time.Millisecond)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Part: 0, Start: 0, Duration: time.Second, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Part: 1, Start: time.Second, Duration: time.Second, Data: []byte("b")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Start: 0, Duration: 2 * time.Second, Data: []byte("ab")}) + + // The next parent starts before the completion hold expires. The previous + // parent must become complete immediately so the playlist has no gap. + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 4, Part: 0, Start: 2 * time.Second, Duration: time.Second, Data: []byte("c")}) + + playlist := w.Playlist("p", "v", completionHoldURI, completionHoldSegment, "init.mp4", nil) + if !strings.Contains(playlist, "#EXTINF:2.000000,") || !strings.Contains(playlist, "4/0.m4s") { + t.Fatalf("playlist lost the held parent when the next parent started:\n%s", playlist) + } + if strings.Contains(playlist, "3/1.m4s") { + t.Fatalf("completed parent still advertised a part:\n%s", playlist) + } +} + +func TestWindowCompletionHoldEvictsAfterTimer(t *testing.T) { + w := NewWindow(WithMaxBytes(2), WithSegmentCompletionDelay(20*time.Millisecond)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Data: []byte("abc")}) + + deadline := time.Now().Add(time.Second) + for time.Now().Before(deadline) && w.Bytes() != 0 { + time.Sleep(5 * time.Millisecond) + } + if got := w.Bytes(); got != 0 { + t.Fatalf("bytes after delayed eviction = %d, want 0", got) + } +} + +func TestWindowCompletionHoldCanceledByPresentationReset(t *testing.T) { + w := NewWindow(WithSegmentCompletionDelay(30 * time.Millisecond)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p1", Track: "v", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p1", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p1", Track: "v", Generation: 1, MSN: 1, Start: 0, Duration: time.Second, Data: []byte("a")}) + + observeEvent(t, w, Event{Kind: Init, Presentation: "p2", Track: "v", Generation: 1, Data: []byte("init2")}) + time.Sleep(60 * time.Millisecond) + if got := w.Snapshot("p1", "v"); len(got.Segments) != 0 { + t.Fatalf("stale segments survived presentation reset: %+v", got.Segments) + } + + observeEvent(t, w, Event{Kind: Part, Presentation: "p2", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("x")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p2", Track: "v", Generation: 1, MSN: 1, Start: 0, Duration: time.Second, Data: []byte("x")}) + waitForSegmentComplete(t, w, "p2", "v", 1) + if got := w.SegmentData("p2", "v", 1); !bytes.Equal(got, []byte("x")) { + t.Fatalf("segment data after reset = %q", got) + } +} + func TestWindowRejectsStalePresentationAndOutOfOrderParts(t *testing.T) { w := NewWindow() if err := w.Observe(Event{Kind: Init, Presentation: "new", Track: "v", Generation: 1, Data: []byte("init")}); err != nil { @@ -75,6 +192,12 @@ func TestWindowEvictsBySegmentsAndBytesAndWakesWaiters(t *testing.T) { if err := w.Observe(Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 2, Part: 0, Data: []byte("bb")}); err != nil { t.Fatal(err) } + if err := w.Observe(Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 1}); err != nil { + t.Fatal(err) + } + if err := w.Observe(Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 2}); err != nil { + t.Fatal(err) + } if err := w.Observe(Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 3, Part: 0, Data: []byte("cc")}); err != nil { t.Fatal(err) } @@ -142,12 +265,12 @@ func TestWindowPreloadHintURIBecomesThePublishedPartURI(t *testing.T) { observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("part-0")}) partURI := func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) } segmentURI := func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) } - playlist := w.Playlist("p", "v", partURI, segmentURI, "init.mp4") + playlist := w.Playlist("p", "v", partURI, segmentURI, "init.mp4", nil) if !strings.Contains(playlist, `#EXT-X-PRELOAD-HINT:TYPE=PART,URI="1/1.m4s"`) { t.Fatalf("playlist omitted preload hint:\n%s", playlist) } observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 1, Data: []byte("part-1")}) - playlist = w.Playlist("p", "v", partURI, segmentURI, "init.mp4") + playlist = w.Playlist("p", "v", partURI, segmentURI, "init.mp4", nil) if !strings.Contains(playlist, `#EXT-X-PART:DURATION=0.000000,URI="1/1.m4s"`) { t.Fatalf("published part did not retain hinted URI:\n%s", playlist) } @@ -242,6 +365,7 @@ func TestWindowExactPartWaitReturnsUnavailableAfterEviction(t *testing.T) { w := NewWindow(WithMaxSegments(1)) observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("old")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 1}) observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 2, Part: 0, Data: []byte("new")}) ctx, cancel := context.WithTimeout(context.Background(), time.Second) @@ -271,7 +395,7 @@ func TestPlaylistContainsLLTagsAndOnlyCompleteParentURI(t *testing.T) { if err := w.Observe(Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 4, Part: 0, Duration: 500 * time.Millisecond, Independent: true, Data: []byte("p")}); err != nil { t.Fatal(err) } - playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) for _, want := range []string{"#EXT-X-VERSION:10", "#EXT-X-PART-INF:", "#EXT-X-SERVER-CONTROL:", `URI="4/0.m4s"`, `#EXT-X-PRELOAD-HINT:TYPE=PART,URI="4/1.m4s"`, "INDEPENDENT=YES"} { if !strings.Contains(playlist, want) { t.Errorf("playlist missing %q:\n%s", want, playlist) @@ -285,6 +409,46 @@ func TestPlaylistContainsLLTagsAndOnlyCompleteParentURI(t *testing.T) { } } +func TestPlaylistSegmentsOnlyOmitsLLPartTags(t *testing.T) { + w := NewWindow() + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "audio", Generation: 1, Data: []byte("init")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 4, Part: 0, Duration: time.Second, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 4, Part: 1, Duration: time.Second, Data: []byte("b")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "audio", Generation: 1, MSN: 4, Duration: 2 * time.Second, Data: []byte("ab")}) + + playlist := w.PlaylistSegmentsOnly("p", "audio", func(msn uint64, part uint32) string { + return fmt.Sprintf("%d/%d.m4s", msn, part) + }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) + for _, forbidden := range []string{"#EXT-X-PART-INF:", "#EXT-X-SERVER-CONTROL:", "#EXT-X-PART:", "#EXT-X-PRELOAD-HINT:", "#EXT-X-RENDITION-REPORT:"} { + if strings.Contains(playlist, forbidden) { + t.Errorf("segments-only playlist contains %q:\n%s", forbidden, playlist) + } + } + if !strings.Contains(playlist, "#EXTINF:2.000000,") || !strings.Contains(playlist, "4.m4s") { + t.Fatalf("segments-only playlist omitted complete parent:\n%s", playlist) + } +} + +func TestPlaylistEmitsIndependentSegmentsOnlyWhenAllParentsStartIndependently(t *testing.T) { + w := NewWindow() + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "video", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Part: 0, Duration: time.Second, Independent: true, Data: []byte("key")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Duration: time.Second, Data: []byte("segment")}) + + playlist := w.Playlist("p", "video", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) + if !strings.Contains(playlist, "#EXT-X-INDEPENDENT-SEGMENTS") { + t.Fatalf("playlist omitted independent-segments declaration for an independent parent:\n%s", playlist) + } + + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 2, Part: 0, Duration: time.Second, Independent: false, Data: []byte("delta")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: 2, Duration: time.Second, Data: []byte("segment")}) + + playlist = w.Playlist("p", "video", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) + if strings.Contains(playlist, "#EXT-X-INDEPENDENT-SEGMENTS") { + t.Fatalf("playlist declared independent segments despite a non-independent parent:\n%s", playlist) + } +} + func TestPlaylistIncludesProgramDateTimeForLiveSegments(t *testing.T) { w := NewWindow() programDateTime := time.Date(2026, time.August, 31, 22, 36, 12, 351000000, time.UTC) @@ -303,12 +467,71 @@ func TestPlaylistIncludesProgramDateTimeForLiveSegments(t *testing.T) { }) observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 4, Start: 2 * time.Second, Duration: time.Second, Data: []byte("segment")}) - playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-PROGRAM-DATE-TIME:2026-08-31T22:36:12.351Z") { t.Fatalf("playlist missing program date time:\n%s", playlist) } } +func TestPlaylistsKeepProgramDateTimeAlignedAcrossTrackDurations(t *testing.T) { + w := NewWindow() + base := time.Date(2026, time.August, 31, 22, 36, 12, 0, time.UTC) + for _, track := range []struct { + name string + parentDur []time.Duration + }{ + {name: "video", parentDur: []time.Duration{2 * time.Second, 2 * time.Second}}, + {name: "audio", parentDur: []time.Duration{1984 * time.Millisecond, 2005 * time.Millisecond}}, + } { + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: track.name, Generation: 1}) + var start time.Duration + for msn, duration := range track.parentDur { + observeEvent(t, w, Event{ + Kind: Part, + Presentation: "p", + Track: track.name, + Generation: 1, + MSN: uint64(msn), + Start: start, + Duration: duration, + ProgramDateTime: base.Add(time.Duration(msn) * 2 * time.Second), + Data: []byte("part"), + }) + observeEvent(t, w, Event{ + Kind: SegmentComplete, + Presentation: "p", + Track: track.name, + Generation: 1, + MSN: uint64(msn), + Start: start, + Duration: duration, + Data: []byte("segment"), + }) + start += duration + } + } + + playlistDates := func(playlist string) []string { + var dates []string + for _, line := range strings.Split(playlist, "\n") { + if strings.HasPrefix(line, "#EXT-X-PROGRAM-DATE-TIME:") { + dates = append(dates, strings.TrimPrefix(line, "#EXT-X-PROGRAM-DATE-TIME:")) + } + } + return dates + } + partURI := func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) } + segmentURI := func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) } + videoDates := playlistDates(w.Playlist("p", "video", partURI, segmentURI, "video-init.mp4", nil)) + audioDates := playlistDates(w.Playlist("p", "audio", partURI, segmentURI, "audio-init.mp4", nil)) + if len(videoDates) != 2 || len(audioDates) != 2 { + t.Fatalf("program date time counts = video %v, audio %v", videoDates, audioDates) + } + if videoDates[1] != audioDates[1] { + t.Fatalf("corresponding parent dates diverged: video=%s audio=%s", videoDates[1], audioDates[1]) + } +} + func TestPlaylistOnlyPublishesPartsForOpenParent(t *testing.T) { w := NewWindow() observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) @@ -319,7 +542,7 @@ func TestPlaylistOnlyPublishesPartsForOpenParent(t *testing.T) { } } - playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if strings.Contains(playlist, `URI="1/0.m4s"`) { t.Fatalf("completed parent still has a part:\n%s", playlist) } @@ -344,7 +567,7 @@ func TestPlaylistRendersAfterInitBeforeFirstPart(t *testing.T) { w := NewWindow() observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) - playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-MAP:URI=\"init.mp4\"") { t.Fatalf("playlist omitted init map:\n%s", playlist) } @@ -359,13 +582,22 @@ func TestWindowStoresVideoConfig(t *testing.T) { } } +func TestWindowStoresAudioConfig(t *testing.T) { + w := NewWindow() + w.SetAudioConfig(AudioConfig{Channels: 1}) + + if got := w.AudioConfig(); got != (AudioConfig{Channels: 1}) { + t.Fatalf("audio config = %+v", got) + } +} + func TestPlaylistDurationsFreezeForPresentationAndRoundTargetDuration(t *testing.T) { w := NewWindow() observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Duration: 1050 * time.Millisecond, Data: []byte("part")}) observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Duration: 5400 * time.Millisecond, Data: []byte("segment")}) - playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:6") { t.Fatalf("playlist did not use the conservative parent duration contract:\n%s", playlist) } @@ -378,7 +610,7 @@ func TestPlaylistDurationsFreezeForPresentationAndRoundTargetDuration(t *testing observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 2, Part: 0, Duration: time.Second, Data: []byte("later-part")}) observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 2, Duration: 6 * time.Second, Data: []byte("later-segment")}) - playlist = w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist = w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:6") || !strings.Contains(playlist, "#EXT-X-PART-INF:PART-TARGET=1.100000") { t.Fatalf("playlist durations changed during presentation:\n%s", playlist) } @@ -389,7 +621,7 @@ func TestPlaylistUsesConfiguredPartTargetWithoutGrowing(t *testing.T) { observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Duration: 500 * time.Millisecond, Data: []byte("part")}) - playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", "v", func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-PART-INF:PART-TARGET=1.100000") { t.Fatalf("playlist did not preserve the presentation part target:\n%s", playlist) } @@ -398,13 +630,13 @@ func TestPlaylistUsesConfiguredPartTargetWithoutGrowing(t *testing.T) { func TestPlaylistDurationContractRoundsAndResetsPerPresentation(t *testing.T) { w := NewWindow(WithPlaylistDurations(2500*time.Millisecond, 750*time.Millisecond)) observeEvent(t, w, Event{Kind: Init, Presentation: "first", Track: "v", Generation: 1}) - playlist := w.Playlist("first", "v", func(uint64, uint32) string { return "part.m4s" }, func(uint64) string { return "segment.m4s" }, "init.mp4") + playlist := w.Playlist("first", "v", func(uint64, uint32) string { return "part.m4s" }, func(uint64) string { return "segment.m4s" }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:3") || !strings.Contains(playlist, "#EXT-X-PART-INF:PART-TARGET=0.750000") { t.Fatalf("playlist did not use configured rounded durations:\n%s", playlist) } observeEvent(t, w, Event{Kind: Init, Presentation: "second", Track: "v", Generation: 1}) - playlist = w.Playlist("second", "v", func(uint64, uint32) string { return "part.m4s" }, func(uint64) string { return "segment.m4s" }, "init.mp4") + playlist = w.Playlist("second", "v", func(uint64, uint32) string { return "part.m4s" }, func(uint64) string { return "segment.m4s" }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:3") || !strings.Contains(playlist, "#EXT-X-PART-INF:PART-TARGET=0.750000") { t.Fatalf("new presentation did not retain configured durations:\n%s", playlist) } @@ -420,9 +652,115 @@ func TestPlaylistUsesOneTargetDurationAcrossTracks(t *testing.T) { observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Duration: 6 * time.Second, Data: []byte("audio")}) for _, track := range []string{"video", "audio"} { - playlist := w.Playlist("p", track, func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4") + playlist := w.Playlist("p", track, func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-TARGETDURATION:6") { t.Errorf("%s playlist did not use the window-wide target duration:\n%s", track, playlist) } } } + +func TestWindowPlaylistRenditionReports(t *testing.T) { + w := NewWindow() + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "video", Generation: 1}) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "audio", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Part: 0, Data: []byte("v")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: 1, Data: []byte("v")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 2, Part: 0, Data: []byte("v")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: 2, Part: 1, Data: []byte("v")}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Part: 0, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "audio", Generation: 1, MSN: 1, Data: []byte("a")}) + + partURI := func(msn uint64, part uint32) string { return fmt.Sprintf("%d/%d.m4s", msn, part) } + segmentURI := func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) } + renditionURI := func(track string) string { return "/r/" + track + "/index.m3u8" } + + // Open audio segment with no parts yet falls back to its completed MSN. + video := w.Playlist("p", "video", partURI, segmentURI, "init.mp4", renditionURI) + if !strings.Contains(video, `#EXT-X-RENDITION-REPORT:URI="/r/audio/index.m3u8",LAST-MSN=1`) || strings.Contains(video, "LAST-PART") { + t.Fatalf("video playlist audio report wrong:\n%s", video) + } + if strings.Contains(video, "/r/video/index.m3u8") { + t.Fatalf("video playlist reports itself:\n%s", video) + } + + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "audio", Generation: 1, MSN: 2, Part: 0, Data: []byte("a")}) + audio := w.Playlist("p", "audio", partURI, segmentURI, "init.mp4", renditionURI) + if !strings.Contains(audio, `#EXT-X-RENDITION-REPORT:URI="/r/video/index.m3u8",LAST-MSN=2,LAST-PART=1`) { + t.Fatalf("audio playlist video report wrong:\n%s", audio) + } + audio = w.Playlist("p", "audio", partURI, segmentURI, "init.mp4", nil) + if strings.Contains(audio, "RENDITION-REPORT") { + t.Fatalf("nil rendition URI must not emit reports:\n%s", audio) + } +} + +func TestWindowEvictionNeverStrandsOpenSegments(t *testing.T) { + // The only segment of the track is still open when its own bytes exceed + // the eviction threshold: it remains in place so its remaining parts and + // completion can still be published. + w := NewWindow(WithMaxBytes(100)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("aaa")}) + + if err := w.Observe(Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 1, Data: []byte("aaa")}); err != nil { + t.Fatalf("open segment was stranded by eviction: %v", err) + } + if err := w.Observe(Event{Kind: SegmentComplete, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Data: []byte("aaaaaa")}); err != nil { + t.Fatalf("segment completion failed after eviction pressure: %v", err) + } +} + +func TestWindowRejectsUnboundedOpenSegmentGrowth(t *testing.T) { + w := NewWindow(WithMaxBytes(2)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) + if err := w.Observe(Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("aaa")}); err != ErrWindowCapacity { + t.Fatalf("open segment overflow error = %v, want %v", err, ErrWindowCapacity) + } + if got := w.Bytes(); got != len("aaa") { + t.Fatalf("bytes after rejected part = %d, want %d", got, len("aaa")) + } +} + +func TestWindowEvictionRetainsHistoryForSmallTracks(t *testing.T) { + // A large track driving the shared byte budget over its limit must not + // drain a small sibling track below the retained history floor. + w := NewWindow(WithMaxBytes(200)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "video", Generation: 1}) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "audio", Generation: 1}) + for msn := uint64(1); msn <= 20; msn++ { + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "video", Generation: 1, MSN: msn, Part: 0, Data: bytes.Repeat([]byte("v"), 20)}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "video", Generation: 1, MSN: msn, Data: bytes.Repeat([]byte("v"), 20)}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "audio", Generation: 1, MSN: msn, Part: 0, Data: []byte("a")}) + observeEvent(t, w, Event{Kind: SegmentComplete, Presentation: "p", Track: "audio", Generation: 1, MSN: msn, Data: []byte("a")}) + } + + got := w.Snapshot("p", "audio").Segments + if len(got) != minRetainedSegments { + t.Fatalf("audio history after eviction = %d segments, want %d", len(got), minRetainedSegments) + } + if first := got[0].MSN; first != 21-minRetainedSegments { + t.Fatalf("audio history starts at MSN %d, want %d", first, 21-minRetainedSegments) + } + if got := w.Snapshot("p", "video").Segments; len(got) == 0 { + t.Fatal("video history was emptied") + } +} + +func TestWindowDiscontinuityReleasesSegmentBytes(t *testing.T) { + w := NewWindow(WithMaxBytes(100)) + observeEvent(t, w, Event{Kind: Init, Presentation: "p", Track: "v", Generation: 1}) + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 1, Part: 0, Data: []byte("part")}) + if got := w.Bytes(); got != len("part") { + t.Fatalf("bytes before discontinuity = %d, want %d", got, len("part")) + } + + observeEvent(t, w, Event{Kind: Discontinuity, Presentation: "p", Track: "v", Generation: 1}) + if got := w.Bytes(); got != 0 { + t.Fatalf("bytes after discontinuity = %d, want 0", got) + } + + observeEvent(t, w, Event{Kind: Part, Presentation: "p", Track: "v", Generation: 1, MSN: 2, Part: 0, Data: []byte("new")}) + if got := w.Bytes(); got != len("new") { + t.Fatalf("bytes after new part = %d, want %d", got, len("new")) + } +} diff --git a/pkg/media/cmaf.go b/pkg/media/cmaf.go index 188aa01a..70494ad4 100644 --- a/pkg/media/cmaf.go +++ b/pkg/media/cmaf.go @@ -3,6 +3,7 @@ package media import ( "bytes" "context" + "errors" "fmt" "sync/atomic" "time" @@ -13,6 +14,13 @@ import ( "stream.place/streamplace/pkg/log" ) +var errUnsupportedCMAFPartLayout = errors.New("unsupported CMAF part layout") + +const ( + llhlsAudioParentMinimum = llhlsParentDuration - 100*time.Millisecond + llhlsAudioParentTolerance = 50 * time.Millisecond +) + // cmafTrackSink translates the buffer-list contract of the GStreamer fMP4 // muxer into application-level LL-HLS events. The first list also carries the // initialization segment. @@ -33,7 +41,27 @@ type cmafTrackSink struct { lastTiming map[uint32]cmafFragmentTiming videoTrackIDs map[uint32]bool hasParent bool + partTarget time.Duration + pendingPart cmafPendingPart samples atomic.Uint64 + // audioOnly marks a track whose samples are all independently decodable + // (AAC) and carry no video, so independence inspection is skipped. + audioOnly bool + // programDateTimeBase anchors the playlist PDT grid. All sinks of one + // presentation share the same base so rendition timelines map to the same + // wall clock; per-parent offsets come from the fragment decode timeline. + programDateTimeBase time.Time + timescale uint32 +} + +type cmafPendingPart struct { + data []byte + start time.Duration + duration time.Duration + independent bool + hasVideo bool + programDateTime time.Time + set bool } func (s *cmafTrackSink) sample(sample *gst.Sample) error { @@ -69,6 +97,22 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { s.videoTrackIDs = videoTrackIDs } } + if s.track == "audio" { + if channels, err := cmafAudioChannels(buffers[initIndex]); err != nil { + if s.ctx != nil { + log.Error(s.ctx, "LL-HLS CMAF audio channel mapping failed", "presentation", s.presentation, "track", s.track, "error", err) + } + } else { + s.window.SetAudioConfig(llhls.AudioConfig{Channels: channels}) + } + } + if timescale, err := cmafTrackTimescale(buffers[initIndex]); err != nil { + if s.ctx != nil { + log.Error(s.ctx, "LL-HLS CMAF track timescale mapping failed", "presentation", s.presentation, "track", s.track, "error", err) + } + } else { + s.timescale = timescale + } if err := s.window.Observe(llhls.Event{ Kind: llhls.Init, Presentation: s.presentation, @@ -89,7 +133,20 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { return nil } first := list.GetBufferAt(uint(mediaStart)) - if first.HasFlags(gst.BufferFlagHeader) && !first.HasFlags(gst.BufferFlagDeltaUnit) && s.hasParent { + chunkDuration := clockDuration(first.Duration()) + if chunkDuration <= 0 { + chunkDuration = s.partDuration + } + // Audio-only muxers can emit a short fragment around a scheduled split. + // Keep that fragment in the current parent when it is close enough to the + // target, but close before accepting another full part that would make the + // parent materially too long. + if s.audioOnly && s.hasParent && s.parentLength >= llhlsAudioParentMinimum && s.parentLength+chunkDuration > llhlsParentDuration+llhlsAudioParentTolerance { + if err := s.completeParent(); err != nil { + return err + } + } + if first.HasFlags(gst.BufferFlagHeader) && !first.HasFlags(gst.BufferFlagDeltaUnit) && s.hasParent && !s.audioOnly { if err := s.completeParent(); err != nil { return err } @@ -102,29 +159,128 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { } fragment.Write(data) } - partStart := s.timelineEnd - chunkDuration := clockDuration(first.Duration()) - if chunkDuration <= 0 { - chunkDuration = s.partDuration + fragmentTiming := cmafFragmentTiming{} + if !s.hasParent && s.timescale > 0 { + timings, err := inspectCMAFFragment(fragment.Bytes()) + if err != nil { + if s.ctx != nil { + log.Error(s.ctx, "LL-HLS CMAF fragment timing unavailable for program date time", "presentation", s.presentation, "track", s.track, "msn", s.nextMSN, "part", s.partIndex, "error", err) + } + } else if len(timings) == 1 { + fragmentTiming = timings[0] + } } + partStart := s.timelineEnd programDateTime := time.Time{} if !s.hasParent { s.parentStart = partStart s.hasParent = true - programDateTime = time.Now().UTC() + if s.programDateTimeBase.IsZero() { + s.programDateTimeBase = time.Now().UTC() + } + programDateTime = s.fragmentProgramDateTime(fragmentTiming, partStart) } s.parentLength += chunkDuration s.timelineEnd += chunkDuration s.parent.Write(fragment.Bytes()) s.inspectTiming(fragment.Bytes()) independent := false - if len(s.videoTrackIDs) > 0 { + hasVideo := false + if s.audioOnly { + // Every sample of the supported audio codecs decodes on its own. + independent = true + } else if len(s.videoTrackIDs) > 0 { var err error independent, err = inspectCMAFFragmentIndependence(fragment.Bytes(), s.videoTrackIDs) if err != nil && s.ctx != nil { log.Error(s.ctx, "LL-HLS CMAF independence inspection failed", "presentation", s.presentation, "track", s.track, "msn", s.nextMSN, "part", s.partIndex, "error", err) independent = false } + hasVideo, err = inspectCMAFFragmentHasVideo(fragment.Bytes(), s.videoTrackIDs) + if err != nil && s.ctx != nil { + log.Error(s.ctx, "LL-HLS CMAF video-track inspection failed", "presentation", s.presentation, "track", s.track, "msn", s.nextMSN, "part", s.partIndex, "error", err) + hasVideo = false + } + } + if err := s.queuePart(cmafPendingPart{ + data: append([]byte(nil), fragment.Bytes()...), + start: partStart, + duration: chunkDuration, + independent: independent, + hasVideo: hasVideo, + programDateTime: programDateTime, + set: true, + }); err != nil { + return err + } + if s.audioOnly && s.parentLength >= llhlsParentDuration { + if err := s.completeParent(); err != nil { + return err + } + } + return nil +} + +func (s *cmafTrackSink) fragmentProgramDateTime(timing cmafFragmentTiming, fallback time.Duration) time.Time { + if s.programDateTimeBase.IsZero() { + s.programDateTimeBase = time.Now().UTC() + } + start := fallback + if s.timescale > 0 { + start = cmafDecodeTimeDuration(timing.DecodeTime, s.timescale) + } + return s.programDateTimeBase.Add(start) +} + +// queuePart keeps the current part unpublished until the following muxer +// chunk proves whether it is terminal. Short audio-only prefixes may be +// merged with the following chunk, but the merged part must remain within the +// advertised PART-TARGET. +func (s *cmafTrackSink) queuePart(next cmafPendingPart) error { + if !next.set || s.partDuration <= 0 || s.partTarget <= 0 { + return s.publishPart(next) + } + if !s.pendingPart.set { + if next.duration >= s.partTargetDuration()*85/100 { + return s.publishPart(next) + } + s.pendingPart = next + return nil + } + target := s.partTargetDuration() + minimum := target * 85 / 100 + if s.pendingPart.duration >= minimum { + if err := s.publishPart(s.pendingPart); err != nil { + return err + } + s.pendingPart = next + return nil + } + if s.pendingPart.duration+next.duration > target { + return fmt.Errorf("%w: short CMAF parts cannot be coalesced within PART-TARGET: prefix=%s next=%s target=%s", errUnsupportedCMAFPartLayout, s.pendingPart.duration, next.duration, target) + } + if !s.pendingPart.hasVideo && next.hasVideo { + s.pendingPart.independent = next.independent + } + s.pendingPart.data = append(s.pendingPart.data, next.data...) + s.pendingPart.duration += next.duration + s.pendingPart.hasVideo = s.pendingPart.hasVideo || next.hasVideo + return nil +} + +func (s *cmafTrackSink) partTargetDuration() time.Duration { + if s.partTarget > 0 { + return s.partTarget + } + return s.partDuration +} + +func (s *cmafTrackSink) publishPart(part cmafPendingPart) error { + if !part.set { + return nil + } + 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{ Kind: llhls.Part, @@ -133,11 +289,11 @@ func (s *cmafTrackSink) sample(sample *gst.Sample) error { Generation: s.generation, MSN: s.nextMSN, Part: s.partIndex, - Start: partStart, - Duration: chunkDuration, - Independent: independent, - ProgramDateTime: programDateTime, - Data: fragment.Bytes(), + Start: part.start, + Duration: part.duration, + Independent: part.independent, + ProgramDateTime: part.programDateTime, + Data: append([]byte(nil), part.data...), }); err != nil { return fmt.Errorf("publish CMAF part: %w", err) } @@ -173,6 +329,10 @@ func (s *cmafTrackSink) completeParent() error { if !s.hasParent { return nil } + if err := s.publishPart(s.pendingPart); err != nil { + return err + } + s.pendingPart = cmafPendingPart{} if err := s.window.Observe(llhls.Event{ Kind: llhls.SegmentComplete, Presentation: s.presentation, @@ -208,18 +368,25 @@ func clockDuration(value gst.ClockTime) time.Duration { func installCMAFSink(ctx context.Context, sink *app.Sink, state *cmafTrackSink) { state.ctx = ctx sink.SetBufferListSupport(true) - sink.SetCallbacks(&app.SinkCallbacks{NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { - sample := sink.PullSample() - if sample == nil { - return gst.FlowEOS - } - if n := state.samples.Add(1); n <= 3 { - log.Log(ctx, "received CMAF buffer list", "presentation", state.presentation, "track", state.track, "sample", n) - } - if err := state.sample(sample); err != nil { - log.Error(ctx, "LL-HLS CMAF output failed", "presentation", state.presentation, "track", state.track, "error", err) - return gst.FlowError - } - return gst.FlowOK - }}) + sink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowEOS + } + if n := state.samples.Add(1); n <= 3 { + log.Log(ctx, "received CMAF buffer list", "presentation", state.presentation, "track", state.track, "sample", n) + } + if err := state.sample(sample); err != nil { + log.Error(ctx, "LL-HLS CMAF output failed", "presentation", state.presentation, "track", state.track, "error", err) + return gst.FlowError + } + return gst.FlowOK + }, + EOSFunc: func(*app.Sink) { + if err := state.completeParent(); err != nil { + log.Error(ctx, "LL-HLS CMAF EOS flush failed", "presentation", state.presentation, "track", state.track, "error", err) + } + }, + }) } diff --git a/pkg/media/cmaf_part_duration_test.go b/pkg/media/cmaf_part_duration_test.go new file mode 100644 index 00000000..80f69cc0 --- /dev/null +++ b/pkg/media/cmaf_part_duration_test.go @@ -0,0 +1,409 @@ +package media + +import ( + "bytes" + "context" + "testing" + "time" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/gstinit" + "stream.place/streamplace/pkg/llhls" +) + +func TestISOFMP4MuxDoesNotPublishShortNonTerminalParts(t *testing.T) { + gstinit.InitGST() + if gst.Find("isofmp4mux") == nil || gst.Find("x264enc") == nil || gst.Find("fdkaacenc") == nil { + t.Skip("static GStreamer build with isofmp4mux, x264enc, and fdkaacenc is required") + } + + pipeline, err := gst.NewPipelineFromString( + "isofmp4mux name=mux fragment-duration=2000000000 chunk-duration=1000000000 send-force-keyunit=false ! appsink name=sink sync=false\n" + + "videotestsrc num-buffers=300 is-live=true pattern=ball ! video/x-raw,width=320,height=240,framerate=30/1 ! x264enc tune=zerolatency speed-preset=ultrafast bframes=0 key-int-max=30 ! h264parse ! video/x-h264,stream-format=avc,alignment=au ! queue ! mux.\n" + + "audiotestsrc num-buffers=480 is-live=true samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc bitrate=128000 ! aacparse ! audio/mpeg,mpegversion=4,stream-format=raw,rate=48000,channels=2 ! queue ! mux.", + ) + if err != nil { + t.Fatal(err) + } + sinkElement, err := pipeline.GetElementByName("sink") + if err != nil { + t.Fatal(err) + } + + state := &cmafTrackSink{ + presentation: "test", + track: "video", + window: llhls.NewWindow(), + partDuration: time.Second, + partTarget: 1100 * time.Millisecond, + } + sink := app.SinkFromElement(sinkElement) + sink.SetBufferListSupport(true) + callbackErr := make(chan error, 1) + done := make(chan struct{}) + sink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowEOS + } + if err := state.sample(sample); err != nil { + select { + case callbackErr <- err: + default: + } + return gst.FlowError + } + return gst.FlowOK + }, + EOSFunc: func(*app.Sink) { close(done) }, + }) + + if err := pipeline.SetState(gst.StatePlaying); err != nil { + t.Fatal(err) + } + select { + case <-done: + case <-time.After(15 * time.Second): + _ = pipeline.SetState(gst.StateNull) + t.Fatal("timed out waiting for irregular-GOP isofmp4mux EOS") + } + _ = pipeline.SetState(gst.StateNull) + select { + case err := <-callbackErr: + t.Fatal(err) + default: + } + + snapshot := state.window.Snapshot("test", "video") + if len(snapshot.Segments) < 2 { + t.Fatalf("expected multiple parent segments, got %d", len(snapshot.Segments)) + } + + const partTarget = 1100 * time.Millisecond + minimumNonTerminal := partTarget * 85 / 100 + var previous []byte + for _, segment := range snapshot.Segments { + for i, part := range segment.Parts { + if part.Duration > partTarget { + t.Fatalf("parent segment %d part %d exceeds PART-TARGET: %s > %s", segment.MSN, part.Index, part.Duration, partTarget) + } + if err := walkCMAFBoxes(part.Data, func(string, []byte) error { return nil }); err != nil { + t.Fatalf("parent segment %d part %d is not parseable CMAF: %v", segment.MSN, part.Index, err) + } + if i+1 < len(segment.Parts) && part.Duration < minimumNonTerminal { + t.Fatalf("parent segment %d part %d is short but non-terminal: %s, minimum %s", segment.MSN, part.Index, part.Duration, minimumNonTerminal) + } + if len(part.Data) == 0 { + t.Fatalf("parent segment %d part %d has no CMAF bytes", segment.MSN, part.Index) + } + if previous != nil && bytes.Equal(previous, part.Data) { + t.Fatalf("parent segment %d part %d aliases the preceding part bytes", segment.MSN, part.Index) + } + previous = append(previous[:0], part.Data...) + } + } +} + +func TestCMAFPartCoalescingKeepsNonIndependentVideoPrefix(t *testing.T) { + const ( + videoTrackID = 2 + partTarget = 1100 * time.Millisecond + minimumPartLength = partTarget * 85 / 100 + ) + prefix := cmafPartDurationTestFragment(videoTrackID, cmafTestNonSyncSampleFlags) + key := cmafPartDurationTestFragment(videoTrackID, cmafTestSyncSampleFlags) + following := cmafPartDurationTestFragment(videoTrackID, cmafTestSyncSampleFlags) + wantMerged := append(append([]byte(nil), prefix...), key...) + + state := &cmafTrackSink{ + presentation: "test", + track: "video", + window: llhls.NewWindow(), + generation: 1, + partDuration: time.Second, + partTarget: partTarget, + videoTrackIDs: map[uint32]bool{videoTrackID: true}, + hasParent: true, + } + if err := state.window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "test", Track: "video", Generation: 1, Data: []byte("init")}); err != nil { + t.Fatal(err) + } + state.parent.Write(prefix) + state.parent.Write(key) + state.parent.Write(following) + state.parentLength = 2 * time.Second + + if err := state.queuePart(cmafPendingPart{ + data: prefix, + start: 0, + duration: 34 * time.Millisecond, + hasVideo: true, + independent: false, + set: true, + }); err != nil { + t.Fatal(err) + } + if got := state.window.Snapshot("test", "video"); len(got.Segments) != 0 { + t.Fatalf("short prefix was published before its successor: %+v", got.Segments) + } + if err := state.queuePart(cmafPendingPart{ + data: key, + start: 34 * time.Millisecond, + duration: time.Second, + hasVideo: true, + independent: true, + set: true, + }); err != nil { + t.Fatal(err) + } + if err := state.queuePart(cmafPendingPart{ + data: following, + start: 1034 * time.Millisecond, + duration: 966 * time.Millisecond, + hasVideo: true, + independent: true, + set: true, + }); err != nil { + t.Fatal(err) + } + + snapshot := state.window.Snapshot("test", "video") + if len(snapshot.Segments) != 1 || len(snapshot.Segments[0].Parts) != 1 { + t.Fatalf("expected only the merged part to be published, got %+v", snapshot.Segments) + } + part := snapshot.Segments[0].Parts[0] + if part.Duration < minimumPartLength || part.Duration > partTarget { + t.Fatalf("merged part duration = %s, want [%s, %s]", part.Duration, minimumPartLength, partTarget) + } + if part.Independent { + t.Fatal("merged part inherited independence from a later keyframe despite its video prefix") + } + if !bytes.Equal(part.Data, wantMerged) { + t.Fatal("merged part bytes do not preserve the original CMAF fragments") + } + if _, err := inspectCMAFFragment(part.Data); err != nil { + t.Fatalf("merged part is not parseable CMAF: %v", err) + } + + prefix[0] ^= 0xff + key[0] ^= 0xff + if !bytes.Equal(snapshot.Segments[0].Parts[0].Data, wantMerged) { + t.Fatal("published merged part changed after source buffers were mutated") + } +} + +func TestCMAFCompleteParentPublishesPendingPart(t *testing.T) { + state := &cmafTrackSink{ + presentation: "test", + track: "video", + window: llhls.NewWindow(), + generation: 1, + partDuration: time.Second, + partTarget: 1100 * time.Millisecond, + hasParent: true, + parentLength: time.Second, + pendingPart: cmafPendingPart{ + data: []byte("pending"), + start: 0, + duration: time.Second, + set: true, + }, + } + if err := state.window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "test", Track: "video", Generation: 1, Data: []byte("init")}); err != nil { + t.Fatal(err) + } + state.parent.WriteString("pending") + + if err := state.completeParent(); err != nil { + t.Fatal(err) + } + snapshot := state.window.Snapshot("test", "video") + if len(snapshot.Segments) != 1 || !snapshot.Segments[0].Complete || len(snapshot.Segments[0].Parts) != 1 { + t.Fatalf("EOS flush did not publish the pending parent: %+v", snapshot.Segments) + } + if !bytes.Equal(snapshot.Segments[0].Parts[0].Data, []byte("pending")) { + t.Fatalf("flushed part data = %q", snapshot.Segments[0].Parts[0].Data) + } +} + +func TestISOFMP4AudioSplitterClosesParents(t *testing.T) { + gstinit.InitGST() + if gst.Find("isofmp4mux") == nil || gst.Find("fdkaacenc") == nil { + t.Skip("static GStreamer build with isofmp4mux and fdkaacenc is required") + } + + pipeline, err := gst.NewPipelineFromString( + "isofmp4mux name=ll_audio_mux manual-split=true fragment-duration=2000000000 chunk-duration=1000000000 ! appsink name=ll_audio_sink sync=false async=false\n" + + "audiotestsrc num-buffers=600 is-live=true timestamp-offset=0 samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! 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.", + ) + if err != nil { + t.Fatal(err) + } + queueElement, err := pipeline.GetElementByName("ll_audio_queue") + if err != nil { + t.Fatal(err) + } + maxSizeTime, err := queueElement.GObject().GetProperty("max-size-time") + if err != nil { + t.Fatal(err) + } + if got, ok := maxSizeTime.(uint64); !ok || got != 0 { + t.Fatalf("manual-split audio queue max-size-time = %v, want 0", maxSizeTime) + } + sinkElement, err := pipeline.GetElementByName("ll_audio_sink") + if err != nil { + t.Fatal(err) + } + async, err := sinkElement.GObject().GetProperty("async") + if err != nil { + t.Fatal(err) + } + if got, ok := async.(bool); !ok || got { + t.Fatalf("CMAF appsink async = %v, want false", async) + } + + state := &cmafTrackSink{ + presentation: "test", + track: "audio", + window: llhls.NewWindow(), + partDuration: time.Second, + partTarget: 1100 * time.Millisecond, + audioOnly: true, + } + sink := app.SinkFromElement(sinkElement) + sink.SetBufferListSupport(true) + callbackErr := make(chan error, 1) + done := make(chan struct{}) + sink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowEOS + } + if err := state.sample(sample); err != nil { + select { + case callbackErr <- err: + default: + } + return gst.FlowError + } + return gst.FlowOK + }, + EOSFunc: func(*app.Sink) { + if err := state.completeParent(); err != nil { + select { + case callbackErr <- err: + default: + } + } + close(done) + }, + }) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + if err := pipeline.SetState(gst.StatePlaying); err != nil { + t.Fatal(err) + } + if err := startLLAudioSplitter(ctx, pipeline); err != nil { + t.Fatal(err) + } + select { + case <-done: + case <-time.After(15 * time.Second): + _ = pipeline.SetState(gst.StateNull) + t.Fatal("timed out waiting for audio splitter isofmp4mux EOS") + } + _ = pipeline.SetState(gst.StateNull) + select { + case err := <-callbackErr: + t.Fatal(err) + default: + } + + snapshot := state.window.Snapshot("test", "audio") + if len(snapshot.Segments) < 4 { + t.Fatalf("audio splitter produced %d parent segments, want at least 4", len(snapshot.Segments)) + } + var have93, have94 bool + for i, segment := range snapshot.Segments { + if i+1 < len(snapshot.Segments) && segment.Duration < 1900*time.Millisecond { + t.Fatalf("audio parent %d is unexpectedly short: %s", segment.MSN, segment.Duration) + } + if i+1 < len(snapshot.Segments) && segment.Duration > 2100*time.Millisecond { + partDurations := make([]time.Duration, 0, len(segment.Parts)) + for _, part := range segment.Parts { + partDurations = append(partDurations, part.Duration) + } + t.Fatalf("audio parent %d is unexpectedly long: %s parts=%v", segment.MSN, segment.Duration, partDurations) + } + if len(segment.Parts) == 0 { + t.Fatalf("audio parent %d has no parts", segment.MSN) + } + var samples uint32 + for _, part := range segment.Parts { + timings, err := inspectCMAFFragment(part.Data) + if err != nil { + t.Fatalf("audio parent %d part %d has invalid timing: %v", segment.MSN, part.Index, err) + } + for _, timing := range timings { + samples += timing.SampleCount + } + } + if i > 0 && i+1 < len(snapshot.Segments) && samples != 93 && samples != 94 { + t.Fatalf("audio parent %d has %d AAC samples, want a frame-aligned 93/94 pattern", segment.MSN, samples) + } + if i > 0 && i+1 < len(snapshot.Segments) { + have93 = have93 || samples == 93 + have94 = have94 || samples == 94 + } + for i, part := range segment.Parts { + if i+1 < len(segment.Parts) && part.Duration < 935*time.Millisecond { + t.Fatalf("audio parent %d part %d is short but non-terminal: %s", segment.MSN, part.Index, part.Duration) + } + } + } + if !have93 || !have94 { + t.Fatalf("audio parents did not alternate AAC frame counts: have93=%v have94=%v", have93, have94) + } +} + +func TestCMAFAudioPartCoalescingKeepsShortPrefixNonTerminal(t *testing.T) { + state := &cmafTrackSink{ + presentation: "test", + track: "audio", + window: llhls.NewWindow(), + generation: 1, + partDuration: time.Second, + partTarget: 1100 * time.Millisecond, + audioOnly: true, + } + if err := state.window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "test", Track: "audio", Generation: 1, Data: []byte("init")}); err != nil { + t.Fatal(err) + } + + for _, part := range []cmafPendingPart{ + {data: []byte("prefix"), duration: 21 * time.Millisecond, set: true}, + {data: []byte("body"), duration: time.Second, set: true}, + {data: []byte("following"), duration: time.Second, set: true}, + } { + if err := state.queuePart(part); err != nil { + t.Fatal(err) + } + } + + snapshot := state.window.Snapshot("test", "audio") + if len(snapshot.Segments) != 1 || len(snapshot.Segments[0].Parts) != 1 { + t.Fatalf("audio short prefix was published separately: %+v", snapshot.Segments) + } + if got := snapshot.Segments[0].Parts[0].Duration; got != 1021*time.Millisecond { + t.Fatalf("coalesced audio part duration = %s, want 1.021s", got) + } +} + +func cmafPartDurationTestFragment(trackID, flags uint32) []byte { + traf := cmafTestIndependenceTraf(trackID, flags, nil, nil) + return append(cmafTestBox("moof", traf), cmafTestBox("mdat", []byte{1, 2, 3, 4})...) +} diff --git a/pkg/media/cmaf_test.go b/pkg/media/cmaf_test.go index ca6bfa52..672edd35 100644 --- a/pkg/media/cmaf_test.go +++ b/pkg/media/cmaf_test.go @@ -26,6 +26,38 @@ func TestIsCMAFInit(t *testing.T) { } } +func TestCMAFPartPublishesFullSizedFirstPartImmediately(t *testing.T) { + state := &cmafTrackSink{ + presentation: "test", + track: "audio", + window: llhls.NewWindow(), + generation: 1, + partDuration: time.Second, + partTarget: 1100 * time.Millisecond, + } + if err := state.window.Observe(llhls.Event{Kind: llhls.Init, Presentation: "test", Track: "audio", Generation: 1, Data: []byte("init")}); err != nil { + t.Fatal(err) + } + + part := cmafPendingPart{ + data: []byte("full-sized-part"), + start: 0, + duration: 939 * time.Millisecond, + set: true, + } + if err := state.queuePart(part); err != nil { + t.Fatal(err) + } + + snapshot := state.window.Snapshot("test", "audio") + if len(snapshot.Segments) != 1 || len(snapshot.Segments[0].Parts) != 1 { + t.Fatalf("full-sized first part was held back: %+v", snapshot.Segments) + } + if got := snapshot.Segments[0].Parts[0].Duration; got != part.duration { + t.Fatalf("published part duration = %s, want %s", got, part.duration) + } +} + func TestCMAFMuxEmitsFragmentedBufferLists(t *testing.T) { gstinit.InitGST() if gst.Find("cmafmux") == nil || gst.Find("x264enc") == nil { @@ -121,7 +153,7 @@ func TestCMAFMuxEmitsFragmentedBufferLists(t *testing.T) { return fmt.Sprintf("%d/%d.m4s", msn, part) }, func(msn uint64) string { return fmt.Sprintf("%d.m4s", msn) - }, "init.mp4") + }, "init.mp4", nil) if !strings.Contains(playlist, "#EXT-X-PROGRAM-DATE-TIME:") { t.Fatalf("CMAF playlist is missing program date time:\n%s", playlist) } diff --git a/pkg/media/cmaf_timing.go b/pkg/media/cmaf_timing.go index 7d491f86..707ff184 100644 --- a/pkg/media/cmaf_timing.go +++ b/pkg/media/cmaf_timing.go @@ -3,6 +3,7 @@ package media import ( "encoding/binary" "fmt" + "time" ) type cmafFragmentTiming struct { @@ -123,6 +124,32 @@ func inspectCMAFFirstVideoSampleIndependent(data []byte, videoTrackID uint32) (b return inspectCMAFFragmentIndependence(data, map[uint32]bool{videoTrackID: true}) } +func inspectCMAFFragmentHasVideo(data []byte, videoTrackIDs map[uint32]bool) (bool, error) { + var foundVideo bool + err := walkCMAFBoxes(data, func(boxType string, payload []byte) error { + if boxType != "moof" { + return nil + } + return walkCMAFBoxes(payload, func(childType string, childPayload []byte) error { + if childType != "traf" { + return nil + } + metadata, err := inspectCMAFTrackFragmentMetadata(childPayload) + if err != nil { + return err + } + if videoTrackIDs[metadata.timing.TrackID] { + foundVideo = true + } + return nil + }) + }) + if err != nil { + return false, err + } + return foundVideo, nil +} + func inspectCMAFFragmentIndependence(data []byte, videoTrackIDs map[uint32]bool) (bool, error) { var foundVideo bool var independent bool @@ -196,6 +223,176 @@ func cmafVideoTrackIDs(data []byte) (map[uint32]bool, error) { return videoTrackIDs, nil } +func cmafTrackTimescale(data []byte) (uint32, error) { + var timescale uint32 + err := walkCMAFBoxes(data, func(boxType string, payload []byte) error { + if boxType != "moov" { + return nil + } + return walkCMAFBoxes(payload, func(childType string, childPayload []byte) error { + if childType != "trak" { + return nil + } + return walkCMAFBoxes(childPayload, func(trackChildType string, trackChildPayload []byte) error { + if trackChildType != "mdia" { + return nil + } + return walkCMAFBoxes(trackChildPayload, func(mediaChildType string, mediaChildPayload []byte) error { + if mediaChildType != "mdhd" || timescale != 0 { + return nil + } + parsed, err := parseCMAFMDHD(mediaChildPayload) + if err == nil { + timescale = parsed + } + return err + }) + }) + }) + }) + if err != nil { + return 0, err + } + if timescale == 0 { + return 0, fmt.Errorf("CMAF init contains no track timescale") + } + return timescale, nil +} + +func cmafAudioChannels(data []byte) (int, error) { + var channels int + var foundAudio bool + err := walkCMAFBoxes(data, func(boxType string, payload []byte) error { + if boxType != "moov" { + return nil + } + return walkCMAFBoxes(payload, func(childType string, childPayload []byte) error { + if childType != "trak" { + return nil + } + trackChannels, found, err := cmafAudioTrackChannels(childPayload) + if err != nil { + return err + } + if !found { + return nil + } + if foundAudio { + return fmt.Errorf("CMAF init contains multiple audio tracks") + } + channels = trackChannels + foundAudio = true + return nil + }) + }) + if err != nil { + return 0, err + } + if !foundAudio { + return 0, fmt.Errorf("CMAF init contains no audio track") + } + return channels, nil +} + +func cmafAudioTrackChannels(data []byte) (int, bool, error) { + var handler string + err := walkCMAFBoxes(data, func(boxType string, payload []byte) error { + if boxType != "mdia" { + return nil + } + return walkCMAFBoxes(payload, func(childType string, childPayload []byte) error { + if childType != "hdlr" { + return nil + } + if len(childPayload) < 12 { + return fmt.Errorf("hdlr is truncated") + } + handler = string(childPayload[8:12]) + return nil + }) + }) + if err != nil { + return 0, false, err + } + if handler != "soun" { + return 0, false, nil + } + + var channels int + err = walkCMAFBoxes(data, func(boxType string, payload []byte) error { + if boxType != "mdia" { + return nil + } + return walkCMAFBoxes(payload, func(childType string, childPayload []byte) error { + if childType != "minf" { + return nil + } + return walkCMAFBoxes(childPayload, func(minfChildType string, minfChildPayload []byte) error { + if minfChildType != "stbl" { + return nil + } + return walkCMAFBoxes(minfChildPayload, func(stblChildType string, stblChildPayload []byte) error { + if stblChildType != "stsd" || channels != 0 { + return nil + } + if len(stblChildPayload) < 8 { + return fmt.Errorf("stsd is truncated") + } + return walkCMAFBoxes(stblChildPayload[8:], func(entryType string, entryPayload []byte) error { + if entryType != "mp4a" || channels != 0 { + return nil + } + if len(entryPayload) < 18 { + return fmt.Errorf("mp4a sample entry is truncated") + } + channels = int(binary.BigEndian.Uint16(entryPayload[16:18])) + if channels == 0 { + return fmt.Errorf("mp4a sample entry has no channels") + } + return nil + }) + }) + }) + }) + }) + if err != nil { + return 0, false, err + } + if channels == 0 { + return 0, false, fmt.Errorf("CMAF audio track contains no mp4a sample entry") + } + return channels, true, nil +} + +func parseCMAFMDHD(payload []byte) (uint32, error) { + if len(payload) < 4 { + return 0, fmt.Errorf("mdhd is truncated") + } + switch payload[0] { + case 0: + if len(payload) < 16 { + return 0, fmt.Errorf("version 0 mdhd is truncated") + } + return binary.BigEndian.Uint32(payload[12:16]), nil + case 1: + if len(payload) < 24 { + return 0, fmt.Errorf("version 1 mdhd is truncated") + } + return binary.BigEndian.Uint32(payload[20:24]), nil + default: + return 0, fmt.Errorf("unsupported mdhd version %d", payload[0]) + } +} + +func cmafDecodeTimeDuration(decodeTime uint64, timescale uint32) time.Duration { + if timescale == 0 { + return 0 + } + wholeSeconds := decodeTime / uint64(timescale) + remainder := decodeTime % uint64(timescale) + return time.Duration(wholeSeconds)*time.Second + time.Duration(remainder)*time.Second/time.Duration(timescale) +} + func parseCMAFTrak(data []byte) (trackID uint32, handler string, err error) { var haveTrackID bool var haveHandler bool diff --git a/pkg/media/cmaf_timing_test.go b/pkg/media/cmaf_timing_test.go index 3c1cd2eb..3b082fd8 100644 --- a/pkg/media/cmaf_timing_test.go +++ b/pkg/media/cmaf_timing_test.go @@ -2,10 +2,77 @@ package media import ( "encoding/binary" + "fmt" "strings" "testing" + "time" ) +func TestCMAFProgramDateTimeUsesFragmentDecodeTime(t *testing.T) { + base := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC) + state := &cmafTrackSink{ + programDateTimeBase: base, + timescale: 48_000, + } + + got := state.fragmentProgramDateTime(cmafFragmentTiming{DecodeTime: 125_000}, time.Second) + want := base.Add(125_000 * time.Second / 48_000) + if !got.Equal(want) { + t.Fatalf("fragment program date time = %s, want %s", got, want) + } +} + +func TestCMAFTrackTimescale(t *testing.T) { + for _, tt := range []struct { + name string + version byte + timescale uint32 + }{ + {name: "version 0", version: 0, timescale: 48_000}, + {name: "version 1", version: 1, timescale: 90_000}, + } { + t.Run(tt.name, func(t *testing.T) { + payloadSize := 16 + if tt.version == 1 { + payloadSize = 24 + } + payload := make([]byte, payloadSize) + payload[0] = tt.version + binary.BigEndian.PutUint32(payload[payloadSize-4:], tt.timescale) + init := cmafTestBox("moov", cmafTestBox("trak", cmafTestBox("mdia", cmafTestBox("mdhd", payload)))) + + got, err := cmafTrackTimescale(init) + if err != nil { + t.Fatal(err) + } + if got != tt.timescale { + t.Fatalf("track timescale = %d, want %d", got, tt.timescale) + } + }) + } +} + +func TestCMAFAudioChannels(t *testing.T) { + for _, channels := range []uint16{1, 2} { + t.Run(fmt.Sprintf("%d channels", channels), func(t *testing.T) { + sampleEntry := make([]byte, 18) + binary.BigEndian.PutUint16(sampleEntry[16:18], channels) + stsdPayload := append(make([]byte, 8), cmafTestBox("mp4a", sampleEntry)...) + hdlrPayload := append(make([]byte, 8), []byte("soun")...) + mediaPayload := append(cmafTestBox("hdlr", hdlrPayload), cmafTestBox("minf", cmafTestBox("stbl", cmafTestBox("stsd", stsdPayload)))...) + init := cmafTestBox("moov", cmafTestBox("trak", cmafTestBox("mdia", mediaPayload))) + + got, err := cmafAudioChannels(init) + if err != nil { + t.Fatal(err) + } + if got != int(channels) { + t.Fatalf("audio channels = %d, want %d", got, channels) + } + }) + } +} + func TestInspectCMAFFragmentReadsTrackTiming(t *testing.T) { tfhd := cmafTestTFHD(2, 9000) tfdt := cmafTestTFDT(1, 0x00000000000f0000) diff --git a/pkg/media/live_window.go b/pkg/media/live_window.go index 9cff51b3..6fb7d334 100644 --- a/pkg/media/live_window.go +++ b/pkg/media/live_window.go @@ -29,12 +29,20 @@ const liveWindowRetention = 30 * time.Second const llhlsWindowSegments = 30 const llhlsWindowBytes = 64 << 20 +// llhlsCompletionHold keeps a finished parent segment appearing open just +// long enough for players blocking on its final part to see that part listed +// and fetch it. The muxer emits the final part and the segment completion in +// the same instant; without the hold the completion wins the playlist update, +// the part request rolls over, and players re-fetch the whole segment. +// Safari audibly replays the segment's audio when that fallback fires. +const llhlsCompletionHold = 300 * time.Millisecond + func (mm *MediaManager) llWindow(did string) *llhls.Window { mm.llWindowsMut.Lock() defer mm.llWindowsMut.Unlock() w := mm.llWindows[did] if w == nil { - w = llhls.NewWindow(llhls.WithMaxSegments(llhlsWindowSegments), llhls.WithMaxBytes(llhlsWindowBytes)) + w = llhls.NewWindow(llhls.WithMaxSegments(llhlsWindowSegments), llhls.WithMaxBytes(llhlsWindowBytes), llhls.WithSegmentCompletionDelay(llhlsCompletionHold)) mm.llWindows[did] = w } return w diff --git a/pkg/media/muxl_segment.go b/pkg/media/muxl_segment.go index 112c2cba..5db77940 100644 --- a/pkg/media/muxl_segment.go +++ b/pkg/media/muxl_segment.go @@ -87,6 +87,9 @@ func muxlSignSegmentElem(ctx context.Context, cli *config.CLI, signStream SignSe appsink, err := gst.NewElementWithProperties("appsink", map[string]any{ "name": "muxl-appsink", "sync": false, + // This sink is part of the live ingest graph. It must not hold the + // pipeline in PAUSED while waiting for a preroll sample. + "async": false, }) if err != nil { return nil, nil, fmt.Errorf("failed to create appsink element: %w", err) diff --git a/pkg/media/rtmp_ingest.go b/pkg/media/rtmp_ingest.go index d2f836dd..065e9b3b 100644 --- a/pkg/media/rtmp_ingest.go +++ b/pkg/media/rtmp_ingest.go @@ -72,12 +72,19 @@ func (mm *MediaManager) RTMPIngest(ctx context.Context, rtmpURL string, ms Media llEnabled = false } if llEnabled { - // Keep audio and video in one CMAF output so both tracks are published - // atomically with one set of parent and part boundaries. + // Audio and video are muxed into separate CMAF renditions: AVPlayer's + // low-latency part stitcher mis-times muxed AAC against the video part + // grid (AAC frames do not divide evenly into it), so audio gets its + // own parts cut on AAC frame boundaries. The renditions are tied + // together by a shared program-date-time anchor. pipelineSlice = append(pipelineSlice, - fmt.Sprintf("isofmp4mux name=ll_video_mux fragment-duration=%d chunk-duration=%d ! appsink name=ll_video_sink sync=false", llhlsParentDuration, llhlsPartDuration), + 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", - "audio_tee. ! queue ! ll_video_mux.", + // The manual-split mux waits for a future AAC buffer to carry each + // split marker. The default one-second queue time limit can fill + // before that buffer arrives and deadlock the muxer. + "audio_tee. ! queue name=ll_audio_queue max-size-time=0 ! ll_audio_mux.", "audio_tee. ! queue name=audio_signer_queue", "demux.video ! queue ! h264parse name=parse ! video/x-h264,stream-format=avc,alignment=au ! tee name=video_tee", "video_tee. ! queue ! ll_video_mux.", @@ -158,6 +165,11 @@ func (mm *MediaManager) RTMPIngest(ctx context.Context, rtmpURL string, ms Media log.Error(ctx, "error setting pipeline to null state", "error", err) } }() + if llEnabled { + if err := startLLAudioSplitter(ctx, pipeline); err != nil { + return err + } + } err = <-busErr log.Log(ctx, "RTMP ingest pipeline stopped", "error", err) @@ -189,6 +201,9 @@ func linkElementToPad(source, destination *gst.Element, sinkPadName string) erro } func installCMAFBranch(ctx context.Context, pipeline *gst.Pipeline, window *llhls.Window, presentation string) error { + // Shared wall-clock anchor so the video and audio rendition playlists map + // to the same program date time (players sync renditions through PDT). + programDateTimeBase := time.Now().UTC() videoElement, err := pipeline.GetElementByName("ll_video_sink") if err != nil { return fmt.Errorf("LL-HLS CMAF sink: %w", err) @@ -217,11 +232,115 @@ func installCMAFBranch(ctx context.Context, pipeline *gst.Pipeline, window *llhl }) } installCMAFSink(ctx, app.SinkFromElement(videoElement), &cmafTrackSink{ - presentation: presentation, - track: "video", - window: window, - generation: 1, - partDuration: llhlsPartDuration, + presentation: presentation, + track: "video", + window: window, + generation: 1, + partDuration: llhlsPartDuration, + programDateTimeBase: programDateTimeBase, + // Must match the PART-TARGET advertised by the window: parts below 85% + // of PART-TARGET make AVPlayer reject the playlist (-12642), so the + // sink coalesces short prefixes up to this target before publishing. + partTarget: 1100 * time.Millisecond, + }) + 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, + track: "audio", + window: window, + generation: 1, + partDuration: llhlsPartDuration, + programDateTimeBase: programDateTimeBase, + partTarget: 1100 * time.Millisecond, + audioOnly: true, + }) + return nil +} + +// startLLAudioSplitter installs a serialized split trigger on the audio queue. +// The manual-split muxer cuts between input buffers, so AAC frames are never +// dropped or assigned to the wrong parent when a 2-second boundary falls +// between two 1024-sample frames. The trigger follows buffer PTS rather than +// wall clock time because live ingest can temporarily run faster than real +// time while upstream data is being drained. +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 99bd4d3a..355e0d56 100644 --- a/pkg/media/rtmp_ingest_test.go +++ b/pkg/media/rtmp_ingest_test.go @@ -1,9 +1,14 @@ package media import ( + "context" "testing" + "time" "github.com/bluenviron/gortsplib/v5/pkg/format" + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/gstinit" ) func TestH264VideoConfigUsesSPSMetadata(t *testing.T) { @@ -21,3 +26,108 @@ func TestH264VideoConfigUsesSPSMetadata(t *testing.T) { t.Fatalf("video config = %+v", config) } } + +func TestLLAudioSplitUsesGstClockTimeSignalType(t *testing.T) { + gstinit.InitGST() + if gst.Find("isofmp4mux") == nil { + t.Skip("static GStreamer build with isofmp4mux is required") + } + + pipeline, err := gst.NewPipelineFromString("isofmp4mux name=mux ! fakesink") + if err != nil { + t.Fatal(err) + } + defer func() { + if err := pipeline.SetState(gst.StateNull); err != nil { + t.Logf("set pipeline to NULL: %v", err) + } + }() + mux, err := pipeline.GetElementByName("mux") + if err != nil { + t.Fatal(err) + } + + if err := emitLLAudioSplit(mux, gst.ClockTime(2*time.Second)); err != nil { + t.Fatalf("audio split signal: %v", err) + } +} + +func TestLLAudioSplitterFollowsMediaTimeline(t *testing.T) { + gstinit.InitGST() + + pipeline, err := gst.NewPipelineFromString("appsrc name=src is-live=true format=time caps=audio/x-raw,format=S16LE,layout=interleaved,rate=48000,channels=1 ! queue name=ll_audio_queue max-size-time=0 ! fakesink name=ll_audio_mux sync=false async=false") + if err != nil { + t.Fatal(err) + } + defer func() { + if err := pipeline.SetState(gst.StateNull); err != nil { + t.Logf("set pipeline to NULL: %v", err) + } + }() + + splitEvents := make(chan bool, 3) + mux, err := pipeline.GetElementByName("ll_audio_mux") + if err != nil { + t.Fatal(err) + } + muxSink := mux.GetStaticPad("sink") + if muxSink == nil { + t.Fatal("LL-HLS audio mux test sink pad is missing") + } + muxSink.AddProbe(gst.PadProbeTypeEventDownstream, func(_ *gst.Pad, info *gst.PadProbeInfo) gst.PadProbeReturn { + event := info.GetEvent() + if event == nil || event.Type() != gst.EventTypeCustomDownstream || !event.HasName("FMP4MuxSplitNow") { + return gst.PadProbeOK + } + value, err := event.GetStructure().GetValue("chunk") + if err != nil { + t.Errorf("read audio split event: %v", err) + return gst.PadProbeOK + } + chunk, ok := value.(bool) + if !ok { + t.Errorf("audio split event chunk = %T(%v), want bool", value, value) + return gst.PadProbeOK + } + splitEvents <- chunk + return gst.PadProbeOK + }) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + if err := startLLAudioSplitter(ctx, pipeline); err != nil { + t.Fatal(err) + } + if err := pipeline.SetState(gst.StatePlaying); err != nil { + t.Fatal(err) + } + + srcElement, err := pipeline.GetElementByName("src") + if err != nil { + t.Fatal(err) + } + src := app.SrcFromElement(srcElement) + if src == nil { + t.Fatal("source element is not appsrc") + } + for i := 0; i <= 6; i++ { + buffer := gst.NewBufferWithSize(1) + buffer.SetPresentationTimestamp(gst.ClockTime(time.Duration(i) * 500 * time.Millisecond)) + buffer.SetDuration(gst.ClockTime(500 * time.Millisecond)) + if result := src.PushBuffer(buffer); result != gst.FlowOK { + t.Fatalf("push audio buffer %d: %s", i, result) + } + } + + want := []bool{true, false, true} + for i, expected := range want { + select { + case got := <-splitEvents: + if got != expected { + t.Fatalf("split event %d chunk = %v, want %v", i, got, expected) + } + case <-time.After(500 * time.Millisecond): + t.Fatalf("timed out waiting for split event %d", i) + } + } +}