From d9481d8bc250f9de8b21577af5a7053ee5f89286 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 14 Nov 2025 13:11:28 -0800 Subject: [PATCH] determinism: who needs average bitrate anyway --- go.mod | 3 +- go.sum | 6 ++-- pkg/cmd/combine.go | 6 ++-- pkg/cmd/split.go | 3 +- pkg/media/segment_converge.go | 68 +++++++++++++++++++++++++++++++++++ pkg/media/segment_split.go | 6 ++-- pkg/media/segmenter.go | 64 ++++++++------------------------- 7 files changed, 98 insertions(+), 58 deletions(-) create mode 100644 pkg/media/segment_converge.go diff --git a/go.mod b/go.mod index b357ee2e..679202d8 100644 --- a/go.mod +++ b/go.mod @@ -13,6 +13,7 @@ replace github.com/bluesky-social/indigo => github.com/streamplace/indigo v0.0.0 require ( firebase.google.com/go/v4 v4.14.1 github.com/99designs/gqlgen v0.17.64 + github.com/Eyevinn/mp4ff v0.50.0 github.com/NYTimes/gziphandler v1.1.1 github.com/ThalesGroup/crypto11 v0.0.0-00010101000000-000000000000 github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d @@ -65,7 +66,6 @@ require ( go.opentelemetry.io/otel v1.36.0 go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.35.0 go.opentelemetry.io/otel/sdk v1.36.0 - go.opentelemetry.io/otel/trace v1.36.0 go.uber.org/goleak v1.3.0 golang.org/x/image v0.30.0 golang.org/x/net v0.43.0 @@ -499,6 +499,7 @@ require ( go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0 // indirect go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.35.0 // indirect go.opentelemetry.io/otel/metric v1.36.0 // indirect + go.opentelemetry.io/otel/trace v1.36.0 // indirect go.opentelemetry.io/proto/otlp v1.5.0 // indirect go.uber.org/atomic v1.11.0 // indirect go.uber.org/automaxprocs v1.6.0 // indirect diff --git a/go.sum b/go.sum index 448c11fb..d5be8d33 100644 --- a/go.sum +++ b/go.sum @@ -104,6 +104,8 @@ github.com/DataDog/zstd v1.4.5 h1:EndNeuB0l9syBZhut0wns3gV1hL8zX8LIu6ZiVHWLIQ= github.com/DataDog/zstd v1.4.5/go.mod h1:1jcaCB/ufaK+sKp1NBhlGmpz41jOoPQ35bpF36t7BBo= github.com/Djarvur/go-err113 v0.0.0-20210108212216-aea10b59be24 h1:sHglBQTwgx+rWPdisA5ynNEsoARbiCBOyGcJM4/OzsM= github.com/Djarvur/go-err113 v0.0.0-20210108212216-aea10b59be24/go.mod h1:4UJr5HIiMZrwgkSPdsjy2uOQExX/WEILpIrO9UPGuXs= +github.com/Eyevinn/mp4ff v0.50.0 h1:vFlsvpQh5Jfz++cuaeTI90vbID5dAabebvvN/l9lom0= +github.com/Eyevinn/mp4ff v0.50.0/go.mod h1:hJNUUqOBryLAzUW9wpCJyw2HaI+TCd2rUPhafoS5lgg= github.com/GaijinEntertainment/go-exhaustruct/v3 v3.3.1 h1:Sz1JIXEcSfhz7fUi7xHnhpIE0thVASYjvosApmHuD2k= github.com/GaijinEntertainment/go-exhaustruct/v3 v3.3.1/go.mod h1:n/LSCXNuIYqVfBlVXyHfMQkZDdp1/mmxfSjADd3z1Zg= github.com/Kagami/go-avif v0.1.0 h1:8GHAGLxCdFfhpd4Zg8j1EqO7rtcQNenxIDerC/uu68w= @@ -463,8 +465,8 @@ github.com/go-sql-driver/mysql v1.8.1/go.mod h1:wEBSXgmK//2ZFJyE+qWnIsVGmvmEKlqw github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY= github.com/go-task/slim-sprig/v3 v3.0.0 h1:sUs3vkvUymDpBKi3qH1YSqBQk9+9D/8M2mN1vB6EwHI= github.com/go-task/slim-sprig/v3 v3.0.0/go.mod h1:W848ghGpv3Qj3dhTPRyJypKRiqCdHZiAzKg9hl15HA8= -github.com/go-test/deep v1.0.8 h1:TDsG77qcSprGbC6vTN8OuXp5g+J+b5Pcguhf7Zt61VM= -github.com/go-test/deep v1.0.8/go.mod h1:5C2ZWiW0ErCdrYzpqxLbTX7MG14M9iiw8DgHncVwcsE= +github.com/go-test/deep v1.1.0 h1:WOcxcdHcvdgThNXjw0t76K42FXTU7HpNQWHpA2HHNlg= +github.com/go-test/deep v1.1.0/go.mod h1:5C2ZWiW0ErCdrYzpqxLbTX7MG14M9iiw8DgHncVwcsE= github.com/go-text/typesetting v0.3.0 h1:OWCgYpp8njoxSRpwrdd1bQOxdjOXDj9Rqart9ML4iF4= github.com/go-text/typesetting v0.3.0/go.mod h1:qjZLkhRgOEYMhU9eHBr3AR4sfnGJvOXNLt8yRAySFuY= github.com/go-text/typesetting-utils v0.0.0-20241103174707-87a29e9e6066 h1:qCuYC+94v2xrb1PoS4NIDe7DGYtLnU2wWiQe9a1B1c0= diff --git a/pkg/cmd/combine.go b/pkg/cmd/combine.go index d2100dfd..960e23ba 100644 --- a/pkg/cmd/combine.go +++ b/pkg/cmd/combine.go @@ -62,19 +62,19 @@ func Combine(ctx context.Context, build *config.BuildFlags, allArgs []string) er if err != nil { return err } - err = CheckCombined(ctx, outFd, *debugDir) + err = CheckCombined(ctx, cli, outFd, *debugDir) if err != nil { return err } return nil } -func CheckCombined(ctx context.Context, inFD io.ReadWriteSeeker, debugDir string) error { +func CheckCombined(ctx context.Context, cli *config.CLI, inFD io.ReadWriteSeeker, debugDir string) error { _, err := inFD.Seek(0, io.SeekStart) if err != nil { return err } - err = media.SplitSegments(ctx, inFD, func(fname string) media.ReadWriteSeekCloser { + err = media.SplitSegments(ctx, cli, inFD, func(fname string) media.ReadWriteSeekCloser { if debugDir == "" { return aqio.NewReadWriteSeeker([]byte{}) } diff --git a/pkg/cmd/split.go b/pkg/cmd/split.go index 987b7ae9..e548c1fc 100644 --- a/pkg/cmd/split.go +++ b/pkg/cmd/split.go @@ -7,6 +7,7 @@ import ( "os" "path/filepath" + "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" ) @@ -24,7 +25,7 @@ func Split(ctx context.Context, inFile, outDir string) error { names := []string{} - err = media.SplitSegments(ctx, inFD, func(fname string) media.ReadWriteSeekCloser { + err = media.SplitSegments(ctx, &config.CLI{}, inFD, func(fname string) media.ReadWriteSeekCloser { fullPath := filepath.Join(outDir, fname) names = append(names, fullPath) log.Log(ctx, "creating segment file", "path", fullPath) diff --git a/pkg/media/segment_converge.go b/pkg/media/segment_converge.go new file mode 100644 index 00000000..f84cbab5 --- /dev/null +++ b/pkg/media/segment_converge.go @@ -0,0 +1,68 @@ +package media + +import ( + "bytes" + "context" + "fmt" + "io" + "os" + "path/filepath" + "slices" + + "github.com/Eyevinn/mp4ff/mp4" + "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/log" +) + +func ConvergeSegment(ctx context.Context, cli *config.CLI, bs []byte, now int64, streamer string) ([]byte, error) { + previousBs := []byte{} + currentBs := bs + i := 0 + for i = 0; i <= MaxSegmentTries; i++ { + if slices.Compare(previousBs, currentBs) == 0 { + break + } + if cli.SegmentDebugDir != "" { + mydir := filepath.Join(cli.SegmentDebugDir, streamer) + err := os.MkdirAll(mydir, 0755) + if err != nil { + return nil, fmt.Errorf("failed to create debug directory: %w", err) + } + aqt := aqtime.FromMillis(now) + outFile := filepath.Join(cli.SegmentDebugDir, fmt.Sprintf("%s-attempt-%03d.mp4", aqt.FileSafeString(), i)) + err = os.WriteFile(outFile, currentBs, 0644) + if err != nil { + return nil, fmt.Errorf("failed to write debug file: %w", err) + } + log.Log(ctx, "wrote debug file", "path", outFile) + } + buf := bytes.Buffer{} + err := CombineSegmentsUnsigned(ctx, []io.ReadSeeker{bytes.NewReader(currentBs)}, &buf) + if err != nil { + return nil, fmt.Errorf("failed to attempt segment convergence: %w", err) + } + previousBs = currentBs + currentBs = buf.Bytes() + mp4file, err := mp4.DecodeFile(bytes.NewReader(currentBs)) + if err != nil { + return nil, fmt.Errorf("failed to decode segment: %w", err) + } + btrt := mp4file.Moov.Trak.Mdia.Minf.Stbl.Stsd.AvcX.Btrt + btrt.AvgBitrate = 0 + btrt.MaxBitrate = 0 + // log.Log(ctx, "btrt", "average bitrate", btrt.AvgBitrate, "max bitrate", btrt.MaxBitrate) + encodedBuf := bytes.Buffer{} + err = mp4file.Encode(&encodedBuf) + if err != nil { + return nil, fmt.Errorf("failed to encode segment: %w", err) + } + currentBs = encodedBuf.Bytes() + } + if slices.Compare(previousBs, currentBs) != 0 { + return nil, fmt.Errorf("failed to converge segment after %d tries", MaxSegmentTries) + } + bs = currentBs + log.Log(ctx, "converged segments", "tries", i, "size", len(bs)) + return currentBs, nil +} diff --git a/pkg/media/segment_split.go b/pkg/media/segment_split.go index 4de2677d..cb6a4406 100644 --- a/pkg/media/segment_split.go +++ b/pkg/media/segment_split.go @@ -11,6 +11,7 @@ import ( "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/aqio" c2patypes "stream.place/streamplace/pkg/c2patypes" + "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/log" ) @@ -95,7 +96,7 @@ func (m *ManySegmentsToSign) Next() *iroh_streamplace.SegmentToSign { } // split a signed concatenated mp4 into its constituent signed segments -func SplitSegments(ctx context.Context, input io.ReadSeeker, cb func(fname string) ReadWriteSeekCloser) error { +func SplitSegments(ctx context.Context, cli *config.CLI, input io.ReadSeeker, cb func(fname string) ReadWriteSeekCloser) error { manifestsStr, err := iroh_streamplace.GetManifests(c2patypes.NewReader(input)) if err != nil { return fmt.Errorf("failed to get manifests: %w", err) @@ -142,6 +143,7 @@ func SplitSegments(ctx context.Context, input io.ReadSeeker, cb func(fname strin } g, ctx := errgroup.WithContext(ctx) unsignedCh := make(chan *SplitSegment) + streamer := manifestList[0].SegmentMetadata.Creator // note: we're passing the input to two places here and need to make sure // they're not running into problems with concurrent seeking. so we use @@ -162,7 +164,7 @@ func SplitSegments(ctx context.Context, input io.ReadSeeker, cb func(fname strin if err != nil { return fmt.Errorf("failed to seek to start: %w", err) } - err = SegmentUnsigned(ctx, input, unsignedCh) + err = SegmentUnsigned(ctx, cli, streamer, input, unsignedCh) if err != nil { return fmt.Errorf("failed to segment file: %w", err) } diff --git a/pkg/media/segmenter.go b/pkg/media/segmenter.go index fd3314a8..bb7febb2 100644 --- a/pkg/media/segmenter.go +++ b/pkg/media/segmenter.go @@ -6,19 +6,17 @@ import ( "fmt" "io" "os" - "path/filepath" - "slices" "strings" "time" "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" - "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" ) // element that takes the input stream, muxes to mp4, and signs the result -func SegmentElem(ctx context.Context, cb func(ctx context.Context, buf []byte, now int64) error) (*gst.Element, error) { +func SegmentElem(ctx context.Context, cli *config.CLI, streamer string, cb func(ctx context.Context, buf []byte, now int64) error) (*gst.Element, error) { // elem, err := gst.NewElement("splitmuxsink name=splitter async-finalize=true sink-factory=appsink muxer-factory=matroskamux max-size-bytes=1") elem, err := gst.NewElementWithProperties("splitmuxsink", map[string]any{ "name": "signer", @@ -121,7 +119,13 @@ func SegmentElem(ctx context.Context, cb func(ctx context.Context, buf []byte, n if previousSegCh != nil { <-previousSegCh } - err := cb(ctx, bs, now) + bs, err := ConvergeSegment(ctx, cli, bs, now, streamer) + if err != nil { + log.Error(ctx, "error converging segment", "error", err) + elem.ErrorMessage(gst.DomainCore, gst.CoreErrorFailed, "Error converging segment", err.Error()) + return + } + err = cb(ctx, bs, now) if err != nil { log.Error(ctx, "error signing segment", "error", err) elem.ErrorMessage(gst.DomainCore, gst.CoreErrorFailed, "Error signing segment", err.Error()) @@ -142,45 +146,7 @@ func SegmentElem(ctx context.Context, cb func(ctx context.Context, buf []byte, n var MaxSegmentTries = 10 func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) (*gst.Element, error) { - return SegmentElem(ctx, func(ctx context.Context, bs []byte, now int64) error { - signedBs, err := ms.SignMP4(ctx, bytes.NewReader(bs), now) - if err != nil { - return fmt.Errorf("error calling SignMP4: %w", err) - } - previousBs := []byte{} - currentBs := signedBs - i := 0 - for i = 0; i <= MaxSegmentTries; i++ { - if slices.Compare(previousBs, currentBs) == 0 { - break - } - if mm.cli.SegmentDebugDir != "" { - mydir := filepath.Join(mm.cli.SegmentDebugDir, ms.Streamer()) - err := os.MkdirAll(mydir, 0755) - if err != nil { - return fmt.Errorf("failed to create debug directory: %w", err) - } - aqt := aqtime.FromMillis(now) - outFile := filepath.Join(mm.cli.SegmentDebugDir, fmt.Sprintf("%s-attempt-%03d.mp4", aqt.FileSafeString(), i)) - err = os.WriteFile(outFile, currentBs, 0644) - if err != nil { - return fmt.Errorf("failed to write debug file: %w", err) - } - log.Log(ctx, "wrote debug file", "path", outFile) - } - buf := bytes.Buffer{} - err := CombineSegmentsUnsigned(ctx, []io.ReadSeeker{bytes.NewReader(currentBs)}, &buf) - if err != nil { - return fmt.Errorf("failed to attempt segment convergence: %w", err) - } - previousBs = currentBs - currentBs = buf.Bytes() - } - if slices.Compare(previousBs, currentBs) != 0 { - return fmt.Errorf("failed to converge segment after %d tries", MaxSegmentTries) - } - bs = currentBs - log.Log(ctx, "converged segments", "tries", i, "size", len(bs)) + return SegmentElem(ctx, mm.cli, ms.Streamer(), func(ctx context.Context, bs []byte, now int64) error { if mm.cli.SmearAudio { smearedBuf := &bytes.Buffer{} err := SmearAudioTimestamps(ctx, bytes.NewReader(bs), smearedBuf) @@ -189,7 +155,7 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) } bs = smearedBuf.Bytes() } - signedBs, err = ms.SignMP4(ctx, bytes.NewReader(bs), now) + signedBs, err := ms.SignMP4(ctx, bytes.NewReader(bs), now) if err != nil { return fmt.Errorf("error calling SignMP4: %w", err) } @@ -202,17 +168,17 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) }) } -func SegmentFileUnsigned(ctx context.Context, input string, ch chan *SplitSegment) error { +func SegmentFileUnsigned(ctx context.Context, cli *config.CLI, streamer string, input string, ch chan *SplitSegment) error { fd, err := os.OpenFile(input, os.O_RDONLY, 0644) log.Log(ctx, "reading file", "file", input) if err != nil { return fmt.Errorf("failed to read file: %w", err) } defer fd.Close() - return SegmentUnsigned(ctx, fd, ch) + return SegmentUnsigned(ctx, cli, streamer, fd, ch) } -func SegmentUnsigned(ctx context.Context, input io.Reader, ch chan *SplitSegment) error { +func SegmentUnsigned(ctx context.Context, cli *config.CLI, streamer string, input io.Reader, ch chan *SplitSegment) error { ctx, cancel := context.WithCancel(ctx) defer cancel() pipelineSlice := []string{ @@ -238,7 +204,7 @@ func SegmentUnsigned(ctx context.Context, input io.Reader, ch chan *SplitSegment return err } - segmenter, err := SegmentElem(ctx, func(ctx context.Context, buf []byte, now int64) error { + segmenter, err := SegmentElem(ctx, cli, streamer, func(ctx context.Context, buf []byte, now int64) error { ch <- &SplitSegment{ Filename: fmt.Sprintf("%d.mp4", now), Data: buf, -- 2.51.2