playback: support serving streams in standard MP4 format (#3213)

* playback: support serving streams in standard MP4 format

* sort samples by DTS

* update readme
This commit is contained in:
Alessandro Ros
2024-04-14 19:29:29 +02:00
committed by GitHub
parent 665e11a376
commit 4157f490fa
10 changed files with 1666 additions and 93 deletions
+8 -1
View File
@@ -1348,7 +1348,7 @@ Where [mypath] is the name of a path. The server will return a list of timespans
The server provides an endpoint for downloading recordings:
```
http://localhost:9996/get?path=[mypath]&start=[start_date]&duration=[duration]
http://localhost:9996/get?path=[mypath]&start=[start_date]&duration=[duration]&format=[format]
```
Where:
@@ -1356,6 +1356,7 @@ Where:
* [mypath] is the path name
* [start_date] is the start date in [RFC3339 format](https://www.utctime.net/)
* [duration] is the maximum duration of the recording in seconds
* [format] (optional) is the output format of the stream. Available values are "fmp4" (default) and "mp4"
All parameters must be [url-encoded](https://www.urlencoder.org/). For instance:
@@ -1371,6 +1372,12 @@ The resulting stream uses the fMP4 format, that is natively compatible with any
</video>
```
The fMP4 format may offer limited compatibility with some players. It's possible to use the standard MP4 format by adding `format=mp4` to a `/get` request:
```
http://localhost:9996/get?path=[mypath]&start=[start_date]&duration=[duration]&format=mp4
```
### Forward streams to other servers
To forward incoming streams to another server, use _FFmpeg_ inside the `runOnReady` parameter:
+83
View File
@@ -0,0 +1,83 @@
package mp4
import (
"io"
"github.com/abema/go-mp4"
)
type mp4Writer struct {
w *mp4.Writer
}
func newMP4Writer(w io.WriteSeeker) *mp4Writer {
return &mp4Writer{
w: mp4.NewWriter(w),
}
}
func (w *mp4Writer) writeBoxStart(box mp4.IImmutableBox) (int, error) {
bi := &mp4.BoxInfo{
Type: box.GetType(),
}
var err error
bi, err = w.w.StartBox(bi)
if err != nil {
return 0, err
}
_, err = mp4.Marshal(w.w, box, mp4.Context{})
if err != nil {
return 0, err
}
return int(bi.Offset), nil
}
func (w *mp4Writer) writeBoxEnd() error {
_, err := w.w.EndBox()
return err
}
func (w *mp4Writer) writeBox(box mp4.IImmutableBox) (int, error) {
off, err := w.writeBoxStart(box)
if err != nil {
return 0, err
}
err = w.writeBoxEnd()
if err != nil {
return 0, err
}
return off, nil
}
func (w *mp4Writer) rewriteBox(off int, box mp4.IImmutableBox) error {
prevOff, err := w.w.Seek(0, io.SeekCurrent)
if err != nil {
return err
}
_, err = w.w.Seek(int64(off), io.SeekStart)
if err != nil {
return err
}
_, err = w.writeBoxStart(box)
if err != nil {
return err
}
err = w.writeBoxEnd()
if err != nil {
return err
}
_, err = w.w.Seek(prevOff, io.SeekStart)
if err != nil {
return err
}
return nil
}
+208
View File
@@ -0,0 +1,208 @@
// Package mp4 contains a MP4 muxer.
package mp4
import (
"io"
"time"
"github.com/abema/go-mp4"
"github.com/bluenviron/mediacommon/pkg/formats/fmp4/seekablebuffer"
)
const (
globalTimescale = 1000
)
func durationMp4ToGo(v int64, timeScale uint32) time.Duration {
timeScale64 := int64(timeScale)
secs := v / timeScale64
dec := v % timeScale64
return time.Duration(secs)*time.Second + time.Duration(dec)*time.Second/time.Duration(timeScale64)
}
// Presentation is timed sequence of video/audio samples.
type Presentation struct {
Tracks []*Track
}
// Marshal encodes a Presentation.
func (p *Presentation) Marshal(w io.Writer) error {
/*
|ftyp|
|moov|
| |mvhd|
| |trak|
| |trak|
| |....|
|mdat|
*/
dataSize, sortedSamples := p.sortSamples()
err := p.marshalFtypAndMoov(w)
if err != nil {
return err
}
return p.marshalMdat(w, dataSize, sortedSamples)
}
func (p *Presentation) sortSamples() (uint32, []*Sample) {
sampleCount := 0
for _, track := range p.Tracks {
sampleCount += len(track.Samples)
}
processedSamples := make([]int, len(p.Tracks))
elapsed := make([]int64, len(p.Tracks))
offset := uint32(0)
sortedSamples := make([]*Sample, sampleCount)
pos := 0
for i, track := range p.Tracks {
elapsed[i] = int64(track.TimeOffset)
}
for {
bestTrack := -1
var bestElapsed time.Duration
for i, track := range p.Tracks {
if processedSamples[i] < len(track.Samples) {
elapsedGo := durationMp4ToGo(elapsed[i], track.TimeScale)
if bestTrack == -1 || elapsedGo < bestElapsed {
bestTrack = i
bestElapsed = elapsedGo
}
}
}
if bestTrack == -1 {
break
}
sample := p.Tracks[bestTrack].Samples[processedSamples[bestTrack]]
sample.offset = offset
processedSamples[bestTrack]++
elapsed[bestTrack] += int64(sample.Duration)
offset += sample.PayloadSize
sortedSamples[pos] = sample
pos++
}
return offset, sortedSamples
}
func (p *Presentation) marshalFtypAndMoov(w io.Writer) error {
var outBuf seekablebuffer.Buffer
mw := newMP4Writer(&outBuf)
_, err := mw.writeBox(&mp4.Ftyp{ // <ftyp/>
MajorBrand: [4]byte{'i', 's', 'o', 'm'},
MinorVersion: 1,
CompatibleBrands: []mp4.CompatibleBrandElem{
{CompatibleBrand: [4]byte{'i', 's', 'o', 'm'}},
{CompatibleBrand: [4]byte{'i', 's', 'o', '2'}},
{CompatibleBrand: [4]byte{'m', 'p', '4', '1'}},
{CompatibleBrand: [4]byte{'m', 'p', '4', '2'}},
},
})
if err != nil {
return err
}
_, err = mw.writeBoxStart(&mp4.Moov{}) // <moov>
if err != nil {
return err
}
mvhd := &mp4.Mvhd{ // <mvhd/>
Timescale: globalTimescale,
Rate: 65536,
Volume: 256,
Matrix: [9]int32{0x00010000, 0, 0, 0, 0x00010000, 0, 0, 0, 0x40000000},
NextTrackID: uint32(len(p.Tracks) + 1),
}
mvhdOffset, err := mw.writeBox(mvhd)
if err != nil {
return err
}
stcos := make([]*mp4.Stco, len(p.Tracks))
stcosOffsets := make([]int, len(p.Tracks))
for i, track := range p.Tracks {
res, err := track.marshal(mw)
if err != nil {
return err
}
stcos[i] = res.stco
stcosOffsets[i] = res.stcoOffset
if res.presentationDuration > mvhd.DurationV0 {
mvhd.DurationV0 = res.presentationDuration
}
}
err = mw.rewriteBox(mvhdOffset, mvhd)
if err != nil {
return err
}
err = mw.writeBoxEnd() // </moov>
if err != nil {
return err
}
moovEndOffset, err := outBuf.Seek(0, io.SeekCurrent)
if err != nil {
return err
}
dataOffset := moovEndOffset + 8
for i := range p.Tracks {
for j := range stcos[i].ChunkOffset {
stcos[i].ChunkOffset[j] += uint32(dataOffset)
}
err = mw.rewriteBox(stcosOffsets[i], stcos[i])
if err != nil {
return err
}
}
_, err = w.Write(outBuf.Bytes())
return err
}
func (p *Presentation) marshalMdat(w io.Writer, dataSize uint32, sortedSamples []*Sample) error {
mdatSize := uint32(8) + dataSize
_, err := w.Write([]byte{byte(mdatSize >> 24), byte(mdatSize >> 16), byte(mdatSize >> 8), byte(mdatSize)})
if err != nil {
return err
}
_, err = w.Write([]byte{'m', 'd', 'a', 't'})
if err != nil {
return err
}
for _, sa := range sortedSamples {
pl, err := sa.GetPayload()
if err != nil {
return err
}
_, err = w.Write(pl)
if err != nil {
return err
}
}
return nil
}
+12
View File
@@ -0,0 +1,12 @@
package mp4
// Sample is a sample of a Track.
type Sample struct {
Duration uint32
PTSOffset int32
IsNonSyncSample bool
PayloadSize uint32
GetPayload func() ([]byte, error)
offset uint32 // filled by sortSamples
}
File diff suppressed because it is too large Load Diff
+7 -1
View File
@@ -5,7 +5,13 @@ import "github.com/bluenviron/mediacommon/pkg/formats/fmp4"
type muxer interface {
writeInit(init *fmp4.Init)
setTrack(trackID int)
writeSample(dts int64, ptsOffset int32, isNonSyncSample bool, payload []byte) error
writeSample(
dts int64,
ptsOffset int32,
isNonSyncSample bool,
payloadSize uint32,
getPayload func() ([]byte, error),
) error
writeFinalDTS(dts int64)
flush() error
}
+22 -10
View File
@@ -9,7 +9,7 @@ import (
)
const (
partSize = 1 * time.Second
partDuration = 1 * time.Second
)
type muxerFMP4Track struct {
@@ -57,12 +57,23 @@ func (w *muxerFMP4) setTrack(trackID int) {
w.curTrack = findTrack(w.tracks, trackID)
}
func (w *muxerFMP4) writeSample(dts int64, ptsOffset int32, isNonSyncSample bool, payload []byte) error {
func (w *muxerFMP4) writeSample(
dts int64,
ptsOffset int32,
isNonSyncSample bool,
_ uint32,
getPayload func() ([]byte, error),
) error {
pl, err := getPayload()
if err != nil {
return err
}
if dts >= 0 {
if w.curTrack.firstDTS < 0 {
w.curTrack.firstDTS = dts
// reset GOP preceding the first frame
// if frame is a IDR, remove previous GOP
if !isNonSyncSample {
w.curTrack.samples = nil
}
@@ -77,29 +88,30 @@ func (w *muxerFMP4) writeSample(dts int64, ptsOffset int32, isNonSyncSample bool
w.curTrack.samples = append(w.curTrack.samples, &fmp4.PartSample{
PTSOffset: ptsOffset,
IsNonSyncSample: isNonSyncSample,
Payload: payload,
Payload: pl,
})
w.curTrack.lastDTS = dts
partSizeMP4 := durationGoToMp4(partSize, w.curTrack.timeScale)
partDurationMP4 := durationGoToMp4(partDuration, w.curTrack.timeScale)
if (w.curTrack.lastDTS - w.curTrack.firstDTS) > partSizeMP4 {
if (w.curTrack.lastDTS - w.curTrack.firstDTS) > partDurationMP4 {
err := w.innerFlush(false)
if err != nil {
return err
}
}
} else {
// store GOP preceding the first frame, with PTSOffset = 0 and Duration = 0
if !isNonSyncSample {
// store GOP of the first frame, and set PTSOffset = 0 and Duration = 0 in each sample
if !isNonSyncSample { // if frame is a IDR, reset GOP
w.curTrack.samples = []*fmp4.PartSample{{
IsNonSyncSample: isNonSyncSample,
Payload: payload,
Payload: pl,
}}
} else {
// append frame to current GOP
w.curTrack.samples = append(w.curTrack.samples, &fmp4.PartSample{
IsNonSyncSample: isNonSyncSample,
Payload: payload,
Payload: pl,
})
}
}
+105
View File
@@ -0,0 +1,105 @@
package playback
import (
"io"
"github.com/bluenviron/mediacommon/pkg/formats/fmp4"
"github.com/bluenviron/mediamtx/internal/playback/mp4"
)
type muxerMP4Track struct {
mp4.Track
lastDTS int64
}
func findTrackMP4(tracks []*muxerMP4Track, id int) *muxerMP4Track {
for _, track := range tracks {
if track.ID == id {
return track
}
}
return nil
}
type muxerMP4 struct {
w io.Writer
tracks []*muxerMP4Track
curTrack *muxerMP4Track
}
func (w *muxerMP4) writeInit(init *fmp4.Init) {
w.tracks = make([]*muxerMP4Track, len(init.Tracks))
for i, track := range init.Tracks {
w.tracks[i] = &muxerMP4Track{
Track: mp4.Track{
ID: track.ID,
TimeScale: track.TimeScale,
Codec: track.Codec,
},
}
}
}
func (w *muxerMP4) setTrack(trackID int) {
w.curTrack = findTrackMP4(w.tracks, trackID)
}
func (w *muxerMP4) writeSample(
dts int64,
ptsOffset int32,
isNonSyncSample bool,
payloadSize uint32,
getPayload func() ([]byte, error),
) error {
// remove GOPs before the GOP of the first frame
if (dts < 0 || (dts >= 0 && w.curTrack.lastDTS < 0)) && !isNonSyncSample {
w.curTrack.Samples = nil
}
if w.curTrack.Samples == nil {
w.curTrack.TimeOffset = int32(dts)
} else {
diff := dts - w.curTrack.lastDTS
if diff < 0 {
diff = 0
}
w.curTrack.Samples[len(w.curTrack.Samples)-1].Duration = uint32(diff)
}
// prevent warning "edit list: 1 Missing key frame while searching for timestamp: 0"
if !isNonSyncSample {
ptsOffset = 0
}
w.curTrack.Samples = append(w.curTrack.Samples, &mp4.Sample{
PTSOffset: ptsOffset,
IsNonSyncSample: isNonSyncSample,
PayloadSize: payloadSize,
GetPayload: getPayload,
})
w.curTrack.lastDTS = dts
return nil
}
func (w *muxerMP4) writeFinalDTS(dts int64) {
diff := dts - w.curTrack.lastDTS
if diff < 0 {
diff = 0
}
w.curTrack.Samples[len(w.curTrack.Samples)-1].Duration = uint32(diff)
}
func (w *muxerMP4) flush() error {
h := mp4.Presentation{
Tracks: make([]*mp4.Track, len(w.tracks)),
}
for i, track := range w.tracks {
h.Tracks[i] = &track.Track
}
return h.Marshal(w.w)
}
+42 -54
View File
@@ -15,8 +15,6 @@ import (
"github.com/gin-gonic/gin"
)
var errStopIteration = errors.New("stop iteration")
type writerWrapper struct {
ctx *gin.Context
written bool
@@ -52,69 +50,52 @@ func seekAndMux(
var firstInit *fmp4.Init
var segmentEnd time.Time
err := func() error {
f, err := os.Open(segments[0].Fpath)
f, err := os.Open(segments[0].Fpath)
if err != nil {
return err
}
defer f.Close()
firstInit, err = segmentFMP4ReadInit(f)
if err != nil {
return err
}
m.writeInit(firstInit)
segmentStartOffset := start.Sub(segments[0].Start)
segmentMaxElapsed, err := segmentFMP4SeekAndMuxParts(f, segmentStartOffset, duration, firstInit, m)
if err != nil {
return err
}
segmentEnd = start.Add(segmentMaxElapsed)
for _, seg := range segments[1:] {
f, err := os.Open(seg.Fpath)
if err != nil {
return err
}
defer f.Close()
firstInit, err = segmentFMP4ReadInit(f)
init, err := segmentFMP4ReadInit(f)
if err != nil {
return err
}
m.writeInit(firstInit)
if !segmentFMP4CanBeConcatenated(firstInit, segmentEnd, init, seg.Start) {
break
}
segmentStartOffset := start.Sub(segments[0].Start)
segmentStartOffset := seg.Start.Sub(start)
segmentMaxElapsed, err := segmentFMP4SeekAndMuxParts(f, segmentStartOffset, duration, firstInit, m)
segmentMaxElapsed, err := segmentFMP4MuxParts(f, segmentStartOffset, duration, firstInit, m)
if err != nil {
return err
}
segmentEnd = start.Add(segmentMaxElapsed)
return nil
}()
if err != nil {
return err
}
for _, seg := range segments[1:] {
err := func() error {
f, err := os.Open(seg.Fpath)
if err != nil {
return err
}
defer f.Close()
init, err := segmentFMP4ReadInit(f)
if err != nil {
return err
}
if !segmentFMP4CanBeConcatenated(firstInit, segmentEnd, init, seg.Start) {
return errStopIteration
}
segmentStartOffset := seg.Start.Sub(start)
segmentMaxElapsed, err := segmentFMP4WriteParts(f, segmentStartOffset, duration, firstInit, m)
if err != nil {
return err
}
segmentEnd = start.Add(segmentMaxElapsed)
return nil
}()
if err != nil {
if errors.Is(err, errStopIteration) {
break
}
return err
}
}
err = m.flush()
@@ -147,8 +128,18 @@ func (p *Server) onGet(ctx *gin.Context) {
return
}
ww := &writerWrapper{ctx: ctx}
var m muxer
format := ctx.Query("format")
if format != "" && format != "fmp4" {
switch format {
case "", "fmp4":
m = &muxerFMP4{w: ww}
case "mp4":
m = &muxerMP4{w: ww}
default:
p.writeError(ctx, http.StatusBadRequest, fmt.Errorf("invalid format: %s", format))
return
}
@@ -169,10 +160,7 @@ func (p *Server) onGet(ctx *gin.Context) {
return
}
ww := &writerWrapper{ctx: ctx}
sw := &muxerFMP4{w: ww}
err = seekAndMux(pathConf.RecordFormat, segments, start, duration, sw)
err = seekAndMux(pathConf.RecordFormat, segments, start, duration, m)
if err != nil {
// user aborted the download
var neterr *net.OpError
+41 -27
View File
@@ -19,6 +19,12 @@ const (
var errTerminated = errors.New("terminated")
type readSeekerAt interface {
io.Reader
io.Seeker
io.ReaderAt
}
func durationGoToMp4(v time.Duration, timeScale uint32) int64 {
timeScale64 := int64(timeScale)
secs := v / time.Second
@@ -337,7 +343,7 @@ func segmentFMP4ReadMaxDuration(
}
func segmentFMP4SeekAndMuxParts(
r io.ReadSeeker,
r readSeekerAt,
segmentStartOffset time.Duration,
duration time.Duration,
init *fmp4.Init,
@@ -394,12 +400,6 @@ func segmentFMP4SeekAndMuxParts(
trun := box.(*mp4.Trun)
dataOffset := moofOffset + uint64(trun.DataOffset)
_, err = r.Seek(int64(dataOffset), io.SeekStart)
if err != nil {
return nil, err
}
muxerDTS := int64(tfdt.BaseMediaDecodeTimeV1) - segmentStartOffsetMP4
atLeastOneSampleWritten := false
@@ -413,23 +413,33 @@ func segmentFMP4SeekAndMuxParts(
atLeastOnePartWritten = true
}
payload := make([]byte, e.SampleSize)
_, err := io.ReadFull(r, payload)
if err != nil {
return nil, err
}
sampleOffset := dataOffset
sampleSize := e.SampleSize
err = m.writeSample(
muxerDTS,
e.SampleCompositionTimeOffsetV1,
(e.SampleFlags&sampleFlagIsNonSyncSample) != 0,
payload,
e.SampleSize,
func() ([]byte, error) {
payload := make([]byte, sampleSize)
n, err := r.ReadAt(payload, int64(sampleOffset))
if err != nil {
return nil, err
}
if n != int(sampleSize) {
return nil, fmt.Errorf("partial read")
}
return payload, nil
},
)
if err != nil {
return nil, err
}
atLeastOneSampleWritten = true
dataOffset += uint64(e.SampleSize)
muxerDTS += int64(e.SampleDuration)
}
@@ -461,8 +471,8 @@ func segmentFMP4SeekAndMuxParts(
return maxMuxerDTS, nil
}
func segmentFMP4WriteParts(
r io.ReadSeeker,
func segmentFMP4MuxParts(
r readSeekerAt,
segmentStartOffset time.Duration,
duration time.Duration,
init *fmp4.Init,
@@ -518,12 +528,6 @@ func segmentFMP4WriteParts(
trun := box.(*mp4.Trun)
dataOffset := moofOffset + uint64(trun.DataOffset)
_, err = r.Seek(int64(dataOffset), io.SeekStart)
if err != nil {
return nil, err
}
muxerDTS := int64(tfdt.BaseMediaDecodeTimeV1) + segmentStartOffsetMP4
atLeastOneSampleWritten := false
@@ -533,23 +537,33 @@ func segmentFMP4WriteParts(
break
}
payload := make([]byte, e.SampleSize)
_, err := io.ReadFull(r, payload)
if err != nil {
return nil, err
}
sampleOffset := dataOffset
sampleSize := e.SampleSize
err = m.writeSample(
muxerDTS,
e.SampleCompositionTimeOffsetV1,
(e.SampleFlags&sampleFlagIsNonSyncSample) != 0,
payload,
e.SampleSize,
func() ([]byte, error) {
payload := make([]byte, sampleSize)
n, err := r.ReadAt(payload, int64(sampleOffset))
if err != nil {
return nil, err
}
if n != int(sampleSize) {
return nil, fmt.Errorf("partial read")
}
return payload, nil
},
)
if err != nil {
return nil, err
}
atLeastOneSampleWritten = true
dataOffset += uint64(e.SampleSize)
muxerDTS += int64(e.SampleDuration)
}