From 21682a34036e62ffb4cdeceff71ada29ddd1a3fd Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 9 Jun 2026 09:14:30 -0700 Subject: [PATCH] vod: swap to flat mp4 storage format --- go.mod | 2 + pkg/muxl/muxl.go | 37 +++++ pkg/spxrpc/place_stream_playback_getvideo.go | 4 +- pkg/vod/finalize_livestream.go | 106 +++---------- pkg/vod/flat_header_test.go | 100 ++++++++++++ pkg/vod/flat_vod.go | 117 ++++++++++++++ pkg/vod/flat_vod_test.go | 158 +++++++++++++++++++ pkg/vod/metafile.go | 45 +++++- pkg/vod/per_segment_metafile_test.go | 104 ++++++++++++ pkg/vod/process.go | 99 +++++++----- pkg/vod/thumbnail_test.go | 2 +- pkg/vod/transfer.go | 17 ++ 12 files changed, 659 insertions(+), 132 deletions(-) create mode 100644 pkg/vod/flat_header_test.go create mode 100644 pkg/vod/flat_vod.go create mode 100644 pkg/vod/flat_vod_test.go create mode 100644 pkg/vod/per_segment_metafile_test.go diff --git a/go.mod b/go.mod index 9a232e56..8b2d6a42 100644 --- a/go.mod +++ b/go.mod @@ -578,3 +578,5 @@ require ( mvdan.cc/unparam v0.0.0-20250301125049-0df0534333a4 // indirect star-tex.org/x/tex v0.7.1 // indirect ) + +replace github.com/streamplace/muxl/go => /home/iameli/code/muxl/go diff --git a/pkg/muxl/muxl.go b/pkg/muxl/muxl.go index 29db6bec..30ec98fa 100644 --- a/pkg/muxl/muxl.go +++ b/pkg/muxl/muxl.go @@ -101,6 +101,43 @@ func RunMuxlWrapInit(ctx context.Context, input io.Reader, output io.Writer) err return eng.WrapInit(ctx, input, output) } +// RunMuxlMetafiles reads a stored MUXL wrapper (canonical fMP4 / flat / bare) +// and writes the payload-free metafile stream — one init then one segment per +// canonical .m4s — as versioned DRISL. Streamplace archives these bytes +// verbatim and feeds them back to RunMuxlSynthesizeFlatHeader. +func RunMuxlMetafiles(ctx context.Context, input io.Reader, output io.Writer) error { + eng, err := getEngine() + if err != nil { + return err + } + return eng.Metafiles(ctx, input, output) +} + +// RunMuxlMetafile returns the payload-free metafile for ONE signed canonical +// segment (muxl metafile --no-init) — the per-fragment archive unit, emitted at +// sign time. The init metafile (catalog) is obtained separately (it's the +// prefix of a RunMuxlMetafiles stream over the first GoP). +func RunMuxlMetafile(ctx context.Context, segment []byte) ([]byte, error) { + eng, err := getEngine() + if err != nil { + return nil, err + } + return eng.Metafile(ctx, segment) +} + +// RunMuxlSynthesizeFlatHeader synthesizes a faststart MP4 header (ftyp + moov + +// mdat-envelope) from a metafile stream (init + N segments). The header's co64 +// offsets are absolute over the [header][body] layout, so serving +// header ++ yields a valid flat MP4 — no base +// offset to pass; muxl owns all the offset math. +func RunMuxlSynthesizeFlatHeader(ctx context.Context, metafiles io.Reader, output io.Writer) error { + eng, err := getEngine() + if err != nil { + return err + } + return eng.SynthesizeFlatHeader(ctx, metafiles, output) +} + // RunMuxlVerify validates the C2PA/S2PA signatures on a signed MUXL wrapper and // returns the per-segment manifest+cert+validation JSON document. func RunMuxlVerify(ctx context.Context, input io.Reader) (string, error) { diff --git a/pkg/spxrpc/place_stream_playback_getvideo.go b/pkg/spxrpc/place_stream_playback_getvideo.go index ccd5e7c1..8e8c2d59 100644 --- a/pkg/spxrpc/place_stream_playback_getvideo.go +++ b/pkg/spxrpc/place_stream_playback_getvideo.go @@ -683,7 +683,9 @@ func mediaPlaylist(meta *vod.Metafile, trackID, ownerDID, sid, cdnURL string, st durSec := float64(seg.DurationTicks) / float64(t.Timescale) lines = append(lines, fmt.Sprintf("#EXTINF:%.6f,", durSec), - fmt.Sprintf("#EXT-X-BYTERANGE:%d@%d", seg.Size, seg.Offset), + // Segment offsets are fragment-relative; the canonical fragments sit + // behind the synthesized flat-MP4 header in the blob, so shift by it. + fmt.Sprintf("#EXT-X-BYTERANGE:%d@%d", seg.Size, seg.Offset+meta.FlatHeaderSize), bURL, ) } diff --git a/pkg/vod/finalize_livestream.go b/pkg/vod/finalize_livestream.go index ef551216..ae2f133c 100644 --- a/pkg/vod/finalize_livestream.go +++ b/pkg/vod/finalize_livestream.go @@ -12,7 +12,6 @@ import ( "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/trace" - "stream.place/streamplace/pkg/bdasl" "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" @@ -78,58 +77,56 @@ func FinalizeLivestreamVOD(ctx context.Context, cli *config.CLI, state *statedb. span.SetAttributes(attribute.Int("object_count", len(keys))) log.Log(ctx, "finalizing livestream VOD", "objects", len(keys)) - // 1. Capture the canonical init header from the first object. We use the - // init `muxl unwrap` itself emits (rather than synthesizing one) so its - // length is exactly what the metafile builder assumes for a leading init — - // keeping the prepended-header blob's offsets self-consistent regardless of - // whether unwrap later passes the header through verbatim or re-derives it. - header, err := captureInitHeader(ctx, store, keys[0]) + // 1. Synthesize the flat-MP4 faststart header from the fragments' metafile. + flatHeader, err := synthFlatHeaderForObjects(ctx, store, keys) if err != nil { - recordErr(span, "capture_init", err) - return "", fmt.Errorf("capture init header: %w", err) + recordErr(span, "synth_header", err) + return "", fmt.Errorf("synthesize flat header: %w", err) } - span.SetAttributes(attribute.Int("header_bytes", len(header))) + span.SetAttributes(attribute.Int("flat_header_bytes", len(flatHeader))) - // 2. One pass over [header]+objects: hash for the CID + build the metafile - // (and per-track init blobs) via unwrap. The content blob is shaped exactly - // like an uploaded VOD, so the metafile offsets are the proven path. - cid, size, metafile, err := hashAndBuildMetafile(ctx, store, header, keys) + // 2. Hash the bare fragments → MUXL CID (the stable, header-independent id) + // and build the fragment-relative HLS metafile + per-track init blobs. + muxlCID, fragSize, metafile, err := hashAndBuildFragmentMetafile(ctx, store, keys) if err != nil { recordErr(span, "build_metafile", err) return "", err } - contentKey := BlobsPrefix + cid + ".mp4" + metafile.FlatHeaderSize = int64(len(flatHeader)) + blobSize := int64(len(flatHeader)) + fragSize + contentKey := BlobsPrefix + muxlCID + ".mp4" span.SetAttributes( - attribute.String("cid", cid), - attribute.Int64("size_bytes", size), + attribute.String("muxl_cid", muxlCID), + attribute.Int64("blob_size", blobSize), attribute.String("content_key", contentKey), ) - // 3. Assemble the content blob, mostly server-side. + // 3. Assemble [flat-header][fragments] at blobs/.mp4 (the flat-header + // is uploaded, the fragments server-side-copied after it). if err := runVODStage(ctx, "concat_assemble", func(ctx context.Context) error { - return assembleContentBlob(ctx, store, header, keys, contentKey) + return assembleContentBlob(ctx, store, flatHeader, keys, contentKey) }); err != nil { recordErr(span, "concat_assemble", err) return "", fmt.Errorf("assemble content blob: %w", err) } if err := runVODStage(ctx, stageMetafile, func(ctx context.Context) error { - return writeMetafile(ctx, store, cid, metafile) + return writeMetafile(ctx, store, muxlCID, metafile) }); err != nil { recordErr(span, stageMetafile, err) return "", fmt.Errorf("write metafile: %w", err) } - // 4. Publish origin + track records and store TrackURIs on the Upload row, - // reusing the upload pipeline's publish path unchanged. + // 4. Publish origin + track records (referencing the MUXL CID) and store + // TrackURIs on the Upload row, reusing the upload publish path unchanged. probe := metafileToVODResult(metafile) if err := runVODStage(ctx, stagePublish, func(ctx context.Context) error { return publishRecords(ctx, publishParams{ cli: cli, state: state, in: Input{UploadID: in.UploadID, RepoDID: in.RepoDID}, - cid: cid, - size: size, + cid: muxlCID, + size: blobSize, mimeType: "video/mp4", probe: probe, signingKey: in.SigningKey, @@ -140,8 +137,8 @@ func FinalizeLivestreamVOD(ctx context.Context, cli *config.CLI, state *statedb. } span.SetStatus(codes.Ok, "") - log.Log(ctx, "livestream VOD finalized", "cid", cid, "size", size, "objects", len(keys), "duration_ms", probe.DurationMS) - return cid, nil + log.Log(ctx, "livestream VOD finalized", "muxlCid", muxlCID, "blobSize", blobSize, "objects", len(keys), "duration_ms", probe.DurationMS) + return muxlCID, nil } // captureInitHeader synthesizes the per-stream init segment (ftyp+moov) from @@ -170,63 +167,6 @@ func captureInitHeader(ctx context.Context, store blob.Store, key string) ([]byt return buf.Bytes(), nil } -// hashAndBuildMetafile streams [header]+objects once, teeing into a bdasl hasher -// (for the content CID) and `muxl unwrap` → metafileBuilder (for the metafile + -// per-track init blobs). Returns the CID, the total blob size, and the -// finalized metafile. -func hashAndBuildMetafile(ctx context.Context, store blob.Store, header []byte, keys []string) (string, int64, *Metafile, error) { - readers := make([]io.Reader, 0, len(keys)+1) - closers := make([]io.Closer, 0, len(keys)) - defer func() { - for _, c := range closers { - _ = c.Close() - } - }() - readers = append(readers, bytes.NewReader(header)) - for _, key := range keys { - r, err := store.Open(ctx, key) - if err != nil { - return "", 0, nil, fmt.Errorf("open object %s: %w", key, err) - } - closers = append(closers, r) - readers = append(readers, io.NewSectionReader(r, 0, r.Size())) - } - - hasher := bdasl.NewWriter() - counter := &countingWriter{} - tee := io.TeeReader(io.MultiReader(readers...), io.MultiWriter(hasher, counter)) - - mb := newMetafileBuilder(ctx, store) - eventCh := make(chan *muxl.MuxlEvent, 16) - producerErr := make(chan error, 1) - go func() { - producerErr <- muxl.RunMuxlUnwrapEvents(ctx, tee, eventCh) - close(eventCh) - }() - - var obsErr error - for ev := range eventCh { - if e := mb.Observe(ev); e != nil && obsErr == nil { - obsErr = e - } - } - if perr := <-producerErr; perr != nil { - return "", 0, nil, fmt.Errorf("muxl unwrap: %w", perr) - } - if obsErr != nil { - return "", 0, nil, fmt.Errorf("metafile build: %w", obsErr) - } - // Unwrap returns nil on context cancellation, which would otherwise let a - // truncated stream yield a partial metafile. Refuse to publish that. - if err := ctx.Err(); err != nil { - return "", 0, nil, err - } - - cid := hasher.CID() - size := counter.load() - return cid, size, mb.Finalize(cid, size), nil -} - // assembleContentBlob writes header ++ objects to contentKey. For an S3 store // this is a near-entirely server-side concat (UploadPartCopy); a non-final // object below S3's 5 MB part floor (only short dev streams) falls back to a diff --git a/pkg/vod/flat_header_test.go b/pkg/vod/flat_header_test.go new file mode 100644 index 00000000..fe453675 --- /dev/null +++ b/pkg/vod/flat_header_test.go @@ -0,0 +1,100 @@ +package vod + +import ( + "bytes" + "context" + "io" + "os" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/bdasl" + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// TestFlatHeaderServing is the end-to-end proof of the flat-MP4 VOD path on the +// streamplace side: from a finalize-shaped canonical blob ([init][signed +// segments]), emit metafiles, synthesize a faststart header, and show that +// serving header ++ via blob.PrefixReader +// reproduces a real flat MP4 byte-for-byte. It also pins down where the flat +// body sits relative to the canonical blob (the bodyOffset the endpoint feeds +// PrefixReader). +// +// Needs gstreamer (warmGST) and the new muxl (Metafiles/SynthesizeFlatHeader), +// so it runs in the cgo container with the local muxl replace. +func TestFlatHeaderServing(t *testing.T) { + warmGST() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + ctx = log.WithLogValues(ctx, "test", "TestFlatHeaderServing") + + fixture, err := os.ReadFile(getFixture("5sec.mp4")) + require.NoError(t, err) + signer, err := newUploadSigner(time.Now()) + require.NoError(t, err) + + // Build a canonical VOD blob exactly like finalize: [init][signed segments]. + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + out := &bytes.Buffer{} + hasher := bdasl.NewWriter() + dst := teeWriter{hasher, out} + mb := newMetafileBuilder(ctx, store) + _, err = streamThroughMuxl(ctx, bytes.NewReader(fixture), int64(len(fixture)), dst, mb, signer.SignerInput) + require.NoError(t, err) + canon := out.Bytes() + meta := mb.Finalize(hasher.CID(), int64(len(canon))) + initLen := minFirstOffset(t, meta) + + // Emit metafiles, then synthesize the flat header from them. + var metas bytes.Buffer + require.NoError(t, muxl.RunMuxlMetafiles(ctx, bytes.NewReader(canon), &metas)) + require.NotZero(t, metas.Len(), "Metafiles produced no output") + var hdr bytes.Buffer + require.NoError(t, muxl.RunMuxlSynthesizeFlatHeader(ctx, bytes.NewReader(metas.Bytes()), &hdr)) + require.NotZero(t, hdr.Len(), "SynthesizeFlatHeader produced no output") + + // Oracle: the real flat MP4 muxl writes from the same canonical blob. + var oracle bytes.Buffer + require.NoError(t, muxl.RunMuxlWrap(ctx, bytes.NewReader(canon), "flat", &oracle)) + + // The synthesized header is the flat MP4's prefix; the remainder is the body. + require.True(t, bytes.HasPrefix(oracle.Bytes(), hdr.Bytes()), + "synthesized header (%d) must be a prefix of the flat MP4 (%d)", hdr.Len(), oracle.Len()) + body := oracle.Bytes()[hdr.Len():] + + // The body must be a contiguous tail range of the canonical blob. + idx := bytes.Index(canon, body) + require.GreaterOrEqual(t, idx, 0, "flat body must be a contiguous range of the canonical blob") + require.Equal(t, len(canon), idx+len(body), "flat body must run to the end of the canonical blob") + t.Logf("flat body offset in canonical blob = %d (initLen=%d, blobSize=%d, headerLen=%d)", + idx, initLen, len(canon), hdr.Len()) + + // Serve header ++ canonical-blob[idx:] via PrefixReader over the STORED blob + // (exactly what the playback endpoint will do) and prove it's byte-identical + // to the real flat MP4. + require.NoError(t, writeBlob(ctx, store, "canon.mp4", canon)) + r, err := store.Open(ctx, "canon.mp4") + require.NoError(t, err) + defer r.Close() + pr := blob.NewPrefixReader(hdr.Bytes(), r, int64(idx), int64(len(body))) + served, err := io.ReadAll(io.NewSectionReader(pr, 0, pr.Size())) + require.NoError(t, err) + require.Equal(t, oracle.Bytes(), served, "PrefixReader must serve the byte-exact flat MP4") +} + +func writeBlob(ctx context.Context, store blob.Store, key string, b []byte) error { + w, err := store.NewWriter(ctx, key, "video/mp4") + if err != nil { + return err + } + defer w.Close() + if _, err := w.Write(b); err != nil { + return err + } + return w.Complete() +} diff --git a/pkg/vod/flat_vod.go b/pkg/vod/flat_vod.go new file mode 100644 index 00000000..73f3bc1a --- /dev/null +++ b/pkg/vod/flat_vod.go @@ -0,0 +1,117 @@ +package vod + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + + "stream.place/streamplace/pkg/bdasl" + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/muxl" +) + +// The flat-MP4 VOD model: the stored blob is [flat-header][canonical fragments], +// content-addressed by the MUXL CID — the BDASL hash of the fragments ALONE, +// not the assembled blob. Keying off the immutable fragments means a better +// flat-header can be re-synthesized and rewritten in place without changing a +// VOD's id (or any media.origin/media.track record), and one blob serves both +// flat-MP4 (the whole file) and HLS (per-track init blobs + byte-ranges into +// the fragments at flatHeaderSize+offset). + +// openFragmentReaders opens the canonical fragment objects as sequential +// readers in order. The caller closes the returned closers. +func openFragmentReaders(ctx context.Context, store blob.Store, keys []string) ([]io.Reader, []io.Closer, error) { + readers := make([]io.Reader, 0, len(keys)) + closers := make([]io.Closer, 0, len(keys)) + for _, key := range keys { + r, err := store.Open(ctx, key) + if err != nil { + for _, c := range closers { + _ = c.Close() + } + return nil, nil, fmt.Errorf("open object %s: %w", key, err) + } + closers = append(closers, r) + readers = append(readers, io.NewSectionReader(r, 0, r.Size())) + } + return readers, closers, nil +} + +func closeAll(closers []io.Closer) { + for _, c := range closers { + _ = c.Close() + } +} + +// hashAndBuildFragmentMetafile streams the bare canonical fragments once, +// teeing into a bdasl hasher (→ the MUXL CID) and muxl unwrap → a +// fragment-relative metafileBuilder (→ the HLS metafile + per-track init +// blobs). The MUXL CID hashes the fragments alone, so it's stable across +// header re-synthesis. Returns the MUXL CID, the fragments' total size, and the +// metafile (with fragment-relative segment offsets). +func hashAndBuildFragmentMetafile(ctx context.Context, store blob.Store, keys []string) (string, int64, *Metafile, error) { + readers, closers, err := openFragmentReaders(ctx, store, keys) + if err != nil { + return "", 0, nil, err + } + defer closeAll(closers) + + hasher := bdasl.NewWriter() + counter := &countingWriter{} + tee := io.TeeReader(io.MultiReader(readers...), io.MultiWriter(hasher, counter)) + + mb := newFragmentMetafileBuilder(ctx, store) + eventCh := make(chan *muxl.MuxlEvent, 16) + producerErr := make(chan error, 1) + go func() { + producerErr <- muxl.RunMuxlUnwrapEvents(ctx, tee, eventCh) + close(eventCh) + }() + + var obsErr error + for ev := range eventCh { + if e := mb.Observe(ev); e != nil && obsErr == nil { + obsErr = e + } + } + if perr := <-producerErr; perr != nil { + return "", 0, nil, fmt.Errorf("muxl unwrap: %w", perr) + } + if obsErr != nil { + return "", 0, nil, fmt.Errorf("metafile build: %w", obsErr) + } + if err := ctx.Err(); err != nil { + return "", 0, nil, err + } + + muxlCID := hasher.CID() + size := counter.load() + return muxlCID, size, mb.Finalize(muxlCID, size), nil +} + +// synthFlatHeaderForObjects emits the metafile stream for the canonical +// fragments and synthesizes the flat-MP4 faststart header from it. muxl owns +// all the co64 offset math over the [header][fragments] layout, so the returned +// header is exactly the bytes to prepend to the fragments. +func synthFlatHeaderForObjects(ctx context.Context, store blob.Store, keys []string) ([]byte, error) { + readers, closers, err := openFragmentReaders(ctx, store, keys) + if err != nil { + return nil, err + } + defer closeAll(closers) + + var metas bytes.Buffer + if err := muxl.RunMuxlMetafiles(ctx, io.MultiReader(readers...), &metas); err != nil { + return nil, fmt.Errorf("emit metafiles: %w", err) + } + var header bytes.Buffer + if err := muxl.RunMuxlSynthesizeFlatHeader(ctx, bytes.NewReader(metas.Bytes()), &header); err != nil { + return nil, fmt.Errorf("synthesize flat header: %w", err) + } + if header.Len() == 0 { + return nil, errors.New("muxl produced an empty flat header") + } + return header.Bytes(), nil +} diff --git a/pkg/vod/flat_vod_test.go b/pkg/vod/flat_vod_test.go new file mode 100644 index 00000000..61187327 --- /dev/null +++ b/pkg/vod/flat_vod_test.go @@ -0,0 +1,158 @@ +package vod + +import ( + "bytes" + "context" + "io" + "os" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/bdasl" + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// TestFinalizeFlatVODAssembly exercises the new finalize model end to end (minus +// statedb/publish): from the bare canonical fragments, it confirms the MUXL CID +// is the fragments' own BDASL hash, the HLS metafile is fragment-relative +// (offsets from 0), and the assembled [flat-header][fragments] blob is a real +// flat MP4 byte-for-byte. Needs gstreamer + the new muxl (local replace). +func TestFinalizeFlatVODAssembly(t *testing.T) { + warmGST() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + ctx = log.WithLogValues(ctx, "test", "TestFinalizeFlatVODAssembly") + + fixture, err := os.ReadFile(getFixture("5sec.mp4")) + require.NoError(t, err) + signer, err := newUploadSigner(time.Now()) + require.NoError(t, err) + + // Produce [init][segments] then strip the init to get the bare fragments + // (what the live recorder stores as .m4s objects). + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + out := &bytes.Buffer{} + h0 := bdasl.NewWriter() + dst := teeWriter{h0, out} + mb0 := newMetafileBuilder(ctx, store) + _, err = streamThroughMuxl(ctx, bytes.NewReader(fixture), int64(len(fixture)), dst, mb0, signer.SignerInput) + require.NoError(t, err) + canon := out.Bytes() + meta0 := mb0.Finalize(h0.CID(), int64(len(canon))) + initLen := minFirstOffset(t, meta0) + fragments := canon[initLen:] + + require.NoError(t, writeBlob(ctx, store, "live/frag.m4s", fragments)) + keys := []string{"live/frag.m4s"} + + // --- the finalize steps --- + flatHeader, err := synthFlatHeaderForObjects(ctx, store, keys) + require.NoError(t, err) + require.NotEmpty(t, flatHeader) + + muxlCID, fragSize, metafile, err := hashAndBuildFragmentMetafile(ctx, store, keys) + require.NoError(t, err) + require.Equal(t, int64(len(fragments)), fragSize, "fragments size") + + // MUXL CID is the BDASL of the bare fragments alone — stable across header + // re-synthesis. + h := bdasl.NewWriter() + _, _ = h.Write(fragments) + require.Equal(t, h.CID(), muxlCID, "MUXL CID must be the BDASL of the bare fragments") + + // HLS metafile is fragment-relative: the earliest segment starts at 0. + require.Equal(t, int64(0), minFirstOffset(t, metafile), "segment offsets must be fragment-relative") + metafile.FlatHeaderSize = int64(len(flatHeader)) + + // Assemble [flat-header][fragments] at blobs/.mp4. + contentKey := BlobsPrefix + muxlCID + ".mp4" + require.NoError(t, assembleContentBlob(ctx, store, flatHeader, keys, contentKey)) + + // The assembled blob is a real flat MP4 (== muxl's own flat wrap), and the + // flat header is its prefix. + var oracle bytes.Buffer + require.NoError(t, muxl.RunMuxlWrap(ctx, bytes.NewReader(fragments), "flat", &oracle)) + cr, err := store.Open(ctx, contentKey) + require.NoError(t, err) + defer cr.Close() + assembled, err := io.ReadAll(io.NewSectionReader(cr, 0, cr.Size())) + require.NoError(t, err) + require.Equal(t, oracle.Bytes(), assembled, "assembled [flat-header][fragments] must be a byte-exact flat MP4") + require.True(t, bytes.HasPrefix(assembled, flatHeader), "flat header must be the blob prefix") + + // HLS reads the same blob: header size + fragment-relative offset lands on + // the segment bytes (== the canonical fragment chunk). + for tid, tr := range metafile.Tracks { + s := tr.Segments[0] + blobOff := metafile.FlatHeaderSize + s.Offset + got := make([]byte, s.Size) + _, err := cr.ReadAt(got, blobOff) + require.NoError(t, err) + require.Equal(t, fragments[s.Offset:s.Offset+s.Size], got, "track %s seg0: HLS byte range must hit the fragment bytes", tid) + } +} + +// TestProcessVODFragmentBuilderPath proves the ProcessVOD path: streaming +// through muxl with a *fragment* metafile builder writes the bare canonical +// fragments to dst (no leading init) — so the same dst hashes to the MUXL CID +// and the synth-header + assemble produces a byte-exact flat MP4. This is the +// inline equivalent of the finalize path's object concatenation, done in a +// single muxl pass with no cross-run determinism assumption. +func TestProcessVODFragmentBuilderPath(t *testing.T) { + warmGST() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + ctx = log.WithLogValues(ctx, "test", "TestProcessVODFragmentBuilderPath") + + fixture, err := os.ReadFile(getFixture("5sec.mp4")) + require.NoError(t, err) + signer, err := newUploadSigner(time.Now()) + require.NoError(t, err) + + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + + // The fragment builder drives streamThroughMuxl to skip the init in dst + // (writeInit=false), exactly as ProcessVOD now does. + out := &bytes.Buffer{} + hasher := bdasl.NewWriter() + mb := newFragmentMetafileBuilder(ctx, store) + _, err = streamThroughMuxl(ctx, bytes.NewReader(fixture), int64(len(fixture)), teeWriter{hasher, out}, mb, signer.SignerInput) + require.NoError(t, err) + fragments := out.Bytes() + require.NotEmpty(t, fragments) + + // dst is bare fragments: the first box is the segment's c2pa 'uuid', NOT + // the init's 'ftyp' — proving the init was kept out of the blob. + require.GreaterOrEqual(t, len(fragments), 8) + require.NotEqual(t, "ftyp", string(fragments[4:8]), "fragments-only blob must not start with the init ftyp box") + + muxlCID := hasher.CID() + require.NoError(t, writeBlob(ctx, store, "stage/frag.m4s", fragments)) + keys := []string{"stage/frag.m4s"} + + flatHeader, err := synthFlatHeaderForObjects(ctx, store, keys) + require.NoError(t, err) + + metafile := mb.Finalize(muxlCID, int64(len(fragments))) + require.Equal(t, int64(0), minFirstOffset(t, metafile), "ProcessVOD metafile must be fragment-relative") + metafile.FlatHeaderSize = int64(len(flatHeader)) + + contentKey := BlobsPrefix + muxlCID + ".mp4" + require.NoError(t, assembleContentBlob(ctx, store, flatHeader, keys, contentKey)) + + var oracle bytes.Buffer + require.NoError(t, muxl.RunMuxlWrap(ctx, bytes.NewReader(fragments), "flat", &oracle)) + cr, err := store.Open(ctx, contentKey) + require.NoError(t, err) + defer cr.Close() + assembled, err := io.ReadAll(io.NewSectionReader(cr, 0, cr.Size())) + require.NoError(t, err) + require.Equal(t, oracle.Bytes(), assembled, "ProcessVOD-style [flat-header][fragments] must be a byte-exact flat MP4") +} diff --git a/pkg/vod/metafile.go b/pkg/vod/metafile.go index 1529152c..96bf4678 100644 --- a/pkg/vod/metafile.go +++ b/pkg/vod/metafile.go @@ -68,6 +68,13 @@ type Metafile struct { BlobCID string `json:"blobCid"` BlobSize int64 `json:"blobSize"` Tracks map[string]MetafileTrack `json:"tracks"` + // FlatHeaderSize is the byte length of the synthesized flat-MP4 header that + // the canonical fragments are stored behind in the blob ([flat-header][ + // fragments]). Segment Offsets are fragment-relative (from 0), so HLS adds + // this to reach the right byte range in the blob; the flat-MP4 path serves + // the whole blob directly. Zero for legacy [init][segments] blobs whose + // offsets are already absolute. + FlatHeaderSize int64 `json:"flatHeaderSize,omitempty"` } // MetafileTrack describes one track within a MUXL container. @@ -124,6 +131,13 @@ type metafileBuilder struct { runningOffset int64 // bytes written to the concatenated output so far seenInit bool + // leadingInitInBlob is true when the blob being measured physically begins + // with the init segment ([init][segments], the legacy/transfer shape), so + // segment offsets start after it. False for the flat-MP4 shape, where the + // stored blob is [flat-header][bare fragments]: segment offsets are + // fragment-relative (from 0) and HLS adds the flat-header size at serve time. + leadingInitInBlob bool + // lastTFDT / tfdtSeen track each track's previous baseMediaDecodeTime so a // backward jump (a concatenated reconnect/restart) can be flagged as a // discontinuity. See MetafileSegment.Discontinuity. @@ -133,15 +147,25 @@ type metafileBuilder struct { func newMetafileBuilder(ctx context.Context, store blob.Store) *metafileBuilder { return &metafileBuilder{ - ctx: ctx, - store: store, - trackInitCIDs: map[string]string{}, - trackSegments: map[string][]MetafileSegment{}, - lastTFDT: map[string]uint64{}, - tfdtSeen: map[string]bool{}, + ctx: ctx, + store: store, + trackInitCIDs: map[string]string{}, + trackSegments: map[string][]MetafileSegment{}, + leadingInitInBlob: true, + lastTFDT: map[string]uint64{}, + tfdtSeen: map[string]bool{}, } } +// newFragmentMetafileBuilder builds a metafile whose segment offsets are +// relative to the bare canonical fragments (from 0), for the flat-MP4 blob +// shape [flat-header][fragments] where the fragments don't start at byte 0. +func newFragmentMetafileBuilder(ctx context.Context, store blob.Store) *metafileBuilder { + b := newMetafileBuilder(ctx, store) + b.leadingInitInBlob = false + return b +} + // Observe processes one MuxlEvent. Order matters — events must arrive // in the same order they're written to the output stream, since that's // how byte offsets get computed. @@ -168,7 +192,14 @@ func (b *metafileBuilder) Observe(ev *muxl.MuxlEvent) error { } b.trackInitCIDs[tid] = cid } - b.runningOffset = int64(len(ev.Data)) + // Advance past the leading init only when it physically prefixes the + // blob. For the flat-MP4 shape the fragments start at 0 (the flat-header + // is added at serve/store time, not measured here). + if b.leadingInitInBlob { + b.runningOffset = int64(len(ev.Data)) + } else { + b.runningOffset = 0 + } case "segment", "signed-segment": // Within a single segment event, per-track byte slices are // concatenated in sorted key order (matching ParseMuxlEvents' diff --git a/pkg/vod/per_segment_metafile_test.go b/pkg/vod/per_segment_metafile_test.go new file mode 100644 index 00000000..9f872727 --- /dev/null +++ b/pkg/vod/per_segment_metafile_test.go @@ -0,0 +1,104 @@ +package vod + +import ( + "bytes" + "context" + "os" + "sort" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/bdasl" + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// TestPerSegmentMetafileAccumulation proves the archival model: emitting one +// RunMuxlMetafile per signed canonical segment (the per-fragment archive unit) +// and concatenating them in canonical byte order reproduces exactly the segment +// portion of the whole-blob RunMuxlMetafiles stream — so the init metafile is +// just the stream prefix, and synthesizing a flat header from +// [init][per-segment metafiles] is byte-identical to synthesizing from the +// whole-blob stream. That's what lets ingest archive a metafile per segment as +// it's signed (no separate Metafiles pass over the fragments). +// +// Needs gstreamer + the new muxl (local replace), so it runs in the cgo container. +func TestPerSegmentMetafileAccumulation(t *testing.T) { + warmGST() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + ctx = log.WithLogValues(ctx, "test", "TestPerSegmentMetafileAccumulation") + + fixture, err := os.ReadFile(getFixture("5sec.mp4")) + require.NoError(t, err) + signer, err := newUploadSigner(time.Now()) + require.NoError(t, err) + + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + out := &bytes.Buffer{} + hasher := bdasl.NewWriter() + dst := teeWriter{hasher, out} + mb := newMetafileBuilder(ctx, store) + _, err = streamThroughMuxl(ctx, bytes.NewReader(fixture), int64(len(fixture)), dst, mb, signer.SignerInput) + require.NoError(t, err) + canon := out.Bytes() + meta := mb.Finalize(hasher.CID(), int64(len(canon))) + + // Whole-blob metafile stream (init + all segments). + var fullStream bytes.Buffer + require.NoError(t, muxl.RunMuxlMetafiles(ctx, bytes.NewReader(canon), &fullStream)) + + // Extract each signed .m4s in canonical byte order (every per-track-per-GoP + // chunk, sorted by offset = the interleave order) and emit its metafile. + type ref struct{ off, size int64 } + var refs []ref + for _, tr := range meta.Tracks { + for _, s := range tr.Segments { + refs = append(refs, ref{s.Offset, s.Size}) + } + } + sort.Slice(refs, func(i, j int) bool { return refs[i].off < refs[j].off }) + + var segStream bytes.Buffer + var firstMF []byte + for i, r := range refs { + m4s := canon[r.off : r.off+r.size] + mf, err := muxl.RunMuxlMetafile(ctx, m4s) + require.NoError(t, err) + require.NotEmpty(t, mf, "per-segment metafile must be non-empty") + if i == 0 { + firstMF = mf + } + segStream.Write(mf) + } + + // Diagnostics: how do the per-segment metafiles relate to the whole stream? + t.Logf("segments=%d fullStream=%d segStream(concat)=%d init≈%d", + len(refs), fullStream.Len(), segStream.Len(), fullStream.Len()-segStream.Len()) + t.Logf("full HasSuffix(segStream)=%v Contains(segStream)=%v", + bytes.HasSuffix(fullStream.Bytes(), segStream.Bytes()), + bytes.Contains(fullStream.Bytes(), segStream.Bytes())) + t.Logf("first per-seg metafile len=%d found in fullStream at idx=%d", + len(firstMF), bytes.Index(fullStream.Bytes(), firstMF)) + + // Whole-blob synth (the proven path) is the oracle header. + var hWhole bytes.Buffer + require.NoError(t, muxl.RunMuxlSynthesizeFlatHeader(ctx, bytes.NewReader(fullStream.Bytes()), &hWhole)) + t.Logf("whole-blob synth header len=%d", hWhole.Len()) + + // Target model: once muxl drops the init/segment distinction (each + // per-segment Metafile becomes self-contained — carries its own catalog), + // the header is synthesized directly from the concatenation of the N + // per-segment metafiles, no separate init. Until that lands, segStream has + // no catalog and the synth errors — skip rather than fail the suite. + var hParts bytes.Buffer + if err := muxl.RunMuxlSynthesizeFlatHeader(ctx, bytes.NewReader(segStream.Bytes()), &hParts); err != nil { + t.Skipf("per-segment metafiles not yet self-contained (pending muxl init-distinction removal): %v", err) + } + require.Equal(t, hWhole.Bytes(), hParts.Bytes(), + "synth from N self-contained per-segment metafiles must match the whole-blob synth") +} diff --git a/pkg/vod/process.go b/pkg/vod/process.go index e6eee133..f26995b2 100644 --- a/pkg/vod/process.go +++ b/pkg/vod/process.go @@ -180,7 +180,11 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, } span.SetAttributes(attribute.String("signing_did", signer.DIDKey)) - metaBuilder := newMetafileBuilder(ctx, store) + // Fragment metafile builder: dst (the staging blob) is the bare canonical + // fragments alone — no leading init — so its hash is the MUXL CID and the + // metafile offsets are fragment-relative. The init is re-synthesized as a + // flat-MP4 faststart header below and prepended at assembly. + metaBuilder := newFragmentMetafileBuilder(ctx, store) probe, err := streamThroughMuxl(ctx, src, size, final, metaBuilder, signer.SignerInput) if err != nil { recordErr(span, stagePipeline, err) @@ -221,25 +225,48 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, } _ = state.SetUploadProgress(ctx, in.UploadID, 90) - finalCID := hasher.CID() - contentKey := BlobsPrefix + finalCID + ".mp4" + // The staging blob is the bare canonical fragments; its BDASL hash is the + // MUXL CID — the stable, header-independent VOD id. Keying the content blob + // off it means a better flat header can later be re-synthesized and + // rewritten in place without changing the VOD's id or any published record. + muxlCID := hasher.CID() + fragSize := counter.load() + contentKey := BlobsPrefix + muxlCID + ".mp4" + + // Synthesize the flat-MP4 faststart header from the fragments' metafile, + // then assemble [flat-header][fragments] at the content key and drop the + // staging blob. muxl owns all co64 offsets over the [header][fragments] + // layout, so the result is a real flat MP4 served whole, with HLS reading + // byte-ranges at flatHeaderSize+offset into the same blob. + flatHeader, err := synthFlatHeaderForObjects(ctx, store, []string{stagingKey}) + if err != nil { + recordErr(span, stageContentAddressCopy, err) + return "", fmt.Errorf("synthesize flat header: %w", err) + } + blobSize := int64(len(flatHeader)) + fragSize span.SetAttributes( - attribute.String("cid", finalCID), + attribute.String("cid", muxlCID), attribute.String("content_key", contentKey), attribute.String("content_url", store.URL(contentKey)), + attribute.Int("flat_header_bytes", len(flatHeader)), + attribute.Int64("blob_size", blobSize), ) if err := runVODStage(ctx, stageContentAddressCopy, func(ctx context.Context) error { - return finalizeMove(ctx, store, stagingKey, contentKey) + if err := assembleContentBlob(ctx, store, flatHeader, []string{stagingKey}, contentKey); err != nil { + return err + } + return store.Delete(ctx, stagingKey) }); err != nil { recordErr(span, stageContentAddressCopy, err) return "", fmt.Errorf("finalize: %w", err) } _ = state.SetUploadProgress(ctx, in.UploadID, 95) - metafile := metaBuilder.Finalize(finalCID, counter.load()) + metafile := metaBuilder.Finalize(muxlCID, fragSize) + metafile.FlatHeaderSize = int64(len(flatHeader)) if err := runVODStage(ctx, stageMetafile, func(ctx context.Context) error { - return writeMetafile(ctx, store, finalCID, metafile) + return writeMetafile(ctx, store, muxlCID, metafile) }); err != nil { recordErr(span, stageMetafile, err) return "", fmt.Errorf("write metafile: %w", err) @@ -254,8 +281,8 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, cli: cli, state: state, in: in, - cid: finalCID, - size: counter.load(), + cid: muxlCID, + size: blobSize, mimeType: "video/mp4", probe: probe, signingKey: signer.DIDKey, @@ -268,13 +295,14 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, spmetrics.VODProcessSuccessesTotal.WithLabelValues(in.Backend).Inc() span.SetStatus(codes.Ok, "") log.Log(ctx, "VOD processed", - "cid", finalCID, + "cid", muxlCID, "url", store.URL(contentKey), "input_size", size, - "output_size", counter.load(), + "frag_size", fragSize, + "output_size", blobSize, "duration_ms", time.Since(startTime).Milliseconds(), ) - return finalCID, nil + return muxlCID, nil } // runVODStage runs a post-pipeline stage with the default stage timeout. @@ -334,28 +362,6 @@ func completeStaging(ctx context.Context, staging blob.Writer, stagingKey string return nil } -// finalizeMove atomically promotes the staged blob at stagingKey to the -// content-addressed contentKey. For S3Store this is a CopyObject + -// DeleteObject; for FileStore it's an os.Rename. -func finalizeMove(ctx context.Context, store blob.Store, stagingKey, contentKey string) error { - ctx, span := vodTracer.Start(ctx, "vod.finalizeMove", trace.WithAttributes( - attribute.String("staging_key", stagingKey), - attribute.String("content_key", contentKey), - )) - defer span.End() - moveStart := time.Now() - if err := store.Move(ctx, stagingKey, contentKey); err != nil { - span.RecordError(err) - span.SetStatus(codes.Error, "move") - return err - } - span.SetAttributes(attribute.Int64("move_duration_ms", time.Since(moveStart).Milliseconds())) - log.Debug(ctx, "moved staging to content-addressed key", - "duration_ms", time.Since(moveStart).Milliseconds(), - ) - return nil -} - // countingWriter is an io.Writer that tallies bytes written via an // atomic counter so a progress goroutine can read it concurrently. type countingWriter struct{ n atomic.Int64 } @@ -540,10 +546,15 @@ func streamThroughMuxl(ctx context.Context, src io.ReaderAt, size int64, dst io. } }() - // Consumer: muxl concatenator output channels -> dst. + // Consumer: muxl concatenator output channels -> dst. The metafile + // builder's offset model dictates the blob layout: leadingInitInBlob + // means dst is [init][segments] (the legacy/transfer shape), otherwise + // dst is the bare fragments alone (the flat-MP4 shape, init synthesized + // separately). A nil builder keeps the legacy behavior. + writeInit := metaBuilder == nil || metaBuilder.leadingInitInBlob consumeDone := make(chan error, 1) go func() { - consumeDone <- consumeConcatTraced(ctx, concat, dst, &initBytes, &segBytes, &initEmits, &segEmits) + consumeDone <- consumeConcatTraced(ctx, concat, dst, writeInit, &initBytes, &segBytes, &initEmits, &segEmits) }() // Optional metafile builder: parallel consumer on the event channel. @@ -616,7 +627,13 @@ func streamThroughMuxl(ctx context.Context, src io.ReaderAt, size int64, dst io. // report them. If the init segment changes mid-stream (multi-input // concatenation), the new init is written too — for VOD with a single // input that doesn't happen, but the loop handles it for free. -func consumeConcatTraced(ctx context.Context, c *muxl.Concatenator, dst io.Writer, initBytes, segBytes, initEmits, segEmits *int64) error { +// +// writeInit controls whether the init segment(s) land in dst. The flat-MP4 +// model wants dst to be the bare canonical fragments alone (the init is +// re-synthesized as a faststart header at finalize time and prepended +// separately); writeInit=false drains and counts the init for tracing but +// keeps it out of the blob — so dst hashes to the MUXL CID of the fragments. +func consumeConcatTraced(ctx context.Context, c *muxl.Concatenator, dst io.Writer, writeInit bool, initBytes, segBytes, initEmits, segEmits *int64) error { initCh, segCh := c.InitCh, c.SegCh start := time.Now() ticker := time.NewTicker(vodProgressLogInterval) @@ -634,9 +651,11 @@ func consumeConcatTraced(ctx context.Context, c *muxl.Concatenator, dst io.Write } *initEmits++ *initBytes += int64(len(init)) - log.Debug(ctx, "muxl init segment", "size", len(init), "emit_n", *initEmits) - if _, err := dst.Write(init); err != nil { - return err + log.Debug(ctx, "muxl init segment", "size", len(init), "emit_n", *initEmits, "in_blob", writeInit) + if writeInit { + if _, err := dst.Write(init); err != nil { + return err + } } case seg, ok := <-segCh: if !ok { diff --git a/pkg/vod/thumbnail_test.go b/pkg/vod/thumbnail_test.go index a98bb1ec..ef6b38ee 100644 --- a/pkg/vod/thumbnail_test.go +++ b/pkg/vod/thumbnail_test.go @@ -45,7 +45,7 @@ func TestGenerateThumbnail(t *testing.T) { meta := mb.Finalize(cid, int64(out.Len())) // generateThumbnail reads the content blob from the store. The real - // ProcessVOD writes it via the staging writer + finalizeMove; this + // ProcessVOD writes it via the staging writer + content assembly; this // test drives streamThroughMuxl directly, so place it by hand. w, err := store.NewWriter(ctx, BlobsPrefix+cid+".mp4", "video/mp4") require.NoError(t, err) diff --git a/pkg/vod/transfer.go b/pkg/vod/transfer.go index 0b207b1b..b9e6a98b 100644 --- a/pkg/vod/transfer.go +++ b/pkg/vod/transfer.go @@ -45,6 +45,23 @@ type TransferResult struct { Downloaded bool `json:"downloaded"` } +// NOTE (flat-MP4 / MUXL-CID migration): TransferVOD has NOT yet been ported +// to the flat-MP4 content model and will fail-safe (CID mismatch, nothing +// stored) on VODs produced by the current finalize/ProcessVOD path. Two +// things break against a [flat-header][fragments] blob keyed by MUXL CID: +// 1. downloadContentBlob hashes the whole downloaded blob and compares to +// contentCID, but contentCID is now BDASL(fragments) alone — the header +// makes the whole-blob hash differ, so verification always mismatches. +// 2. regenerateSidecars runs `muxl unwrap` over the blob, but a flat +// presentation MP4 (ftyp+moov+co64) isn't a MUXL wrapper unwrap can read, +// and it uses the leading-init metafile builder (wrong offset base). +// +// The fix needs the receiver to recover the bare fragments from the blob +// (strip the flat header) before hashing/unwrapping. Since box-walking should +// live in muxl (not here), this awaits either a muxl "strip flat header" / +// "unwrap flat" capability or a wire-protocol change (e.g. transfer the bare +// fragments, or ship FlatHeaderSize alongside). Tracked as a follow-up. +// // TransferVOD pulls a content-addressed VOD blob from a remote Streamplace // node into this node's playback store, regenerates the playback sidecars // (metafile + per-track init segments) locally from the blob, and publishes -- 2.51.2