diff --git a/pkg/director/s3_upload.go b/pkg/director/s3_upload.go index 70a8de1df..a2dad2499 100644 --- a/pkg/director/s3_upload.go +++ b/pkg/director/s3_upload.go @@ -30,7 +30,9 @@ func (ss *StreamSession) s3Upload(ctx context.Context, notif *media.NewSegmentNo return } ss.Go(ctx, func() error { - return ss.s3Uploader.AddSegment(ctx, notif.Data) + // notif.Muxl is the bare canonical segment; it concatenates directly + // (the S3 uploader synthesizes one init and prepends it per object). + return ss.s3Uploader.AddSegment(ctx, notif.Muxl) }) } diff --git a/pkg/muxl/muxl.go b/pkg/muxl/muxl.go index b98939a60..29db6becc 100644 --- a/pkg/muxl/muxl.go +++ b/pkg/muxl/muxl.go @@ -178,10 +178,10 @@ func RunMuxlSignTranscode(ctx context.Context, in TranscodeInput) ([]byte, error // the upstream value can't be built without a ready Engine, and callers expect // construction never to fail (the error surfaces on Close). // -// cat := muxl.NewConcatenator(ctx) +// cat := muxl.NewSigningSegmenter(ctx, signerInput) // go func() { cat.Write(fullFmp4Archive); cat.Close() }() // initSeg := <-cat.InitCh -// for seg := range cat.SegCh { /* append to output */ } +// for seg := range cat.SegCh { /* append signed segments to output */ } type Concatenator struct { // InitCh receives an init segment only when the track configuration // changes. SegCh receives concatenable segment bodies (signed, for the @@ -202,17 +202,6 @@ func (c *Concatenator) Write(data []byte) error { return c.write(data) } // channels are closed by the time it returns. func (c *Concatenator) Close() error { return c.closeFn() } -// NewConcatenator drives Engine.ConcatEvents: dedup/concatenate MUXL fMP4 -// archives into one stream (a new init is emitted only when the catalog -// changes). -func NewConcatenator(ctx context.Context) *Concatenator { - eng, err := getEngine() - if err != nil { - return failedConcatenator(err) - } - return adopt(upstream.NewConcatenator(ctx, eng)) -} - // NewSigningSegmenter drives Engine.SignSegment: each canonical segment is // S2PA-signed in place, so SegCh carries [c2pa-uuid][muxl-uuid][moof][mdat] per // track. Exactly one of in.KeyPEM or in.Sign must be set. diff --git a/pkg/muxl/muxl_sign_test.go b/pkg/muxl/muxl_sign_test.go deleted file mode 100644 index 697d1e602..000000000 --- a/pkg/muxl/muxl_sign_test.go +++ /dev/null @@ -1,73 +0,0 @@ -package muxl - -import ( - "context" - "os" - "path/filepath" - "sync" - "testing" - "time" - - "github.com/stretchr/testify/require" - "stream.place/streamplace/test/remote" -) - -// TestConcatenatorRealSegments drives the wasm Concatenator with real -// signed flat MP4 segments and confirms it -// emits one segment event per input file (minus one held back until -// Close, which reaches flush() in the wasm). -// -// Reproduces the bug where SegCh stayed silent: concat used to treat the -// outer mdat envelope as opaque and drop all the inner moof+mdat -// fragments. Skipped if no segments are available locally. -func TestConcatenatorRealSegments(t *testing.T) { - segPaths := []string{ - remote.RemoteFixture("507b7782c4a6855863a5d4c32cbac3e8fc026c9b635dd174d23c889b030dfc71/2026-05-08T21-55-06-837Z.mp4"), - remote.RemoteFixture("fa42021f9fef60213d801f3521ef725e27c623899ab621a8615fa6107e27a997/2026-05-08T21-55-07-837Z.mp4"), - remote.RemoteFixture("a8b98878338f3946b257e2b68f63157e1fe81a571d5e1954bee093db293c14d3/2026-05-08T21-55-08-838Z.mp4"), - remote.RemoteFixture("2a57daaf03afecc45621ba25794cb87e30fb0a3ae7a18961a6f8045f937d8c4a/2026-05-08T21-55-09-836Z.mp4"), - remote.RemoteFixture("5658128b0c766fb4b20e748ba834df7151c58c0be2f3801e7234d9736a687fc3/2026-05-08T21-55-10-837Z.mp4"), - } - - ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) - defer cancel() - - concat := NewConcatenator(ctx) - - // Spin up consumers BEFORE writing so the producer doesn't block on a - // full SegCh. - var ( - initEvents int - segEvents int - consumeWg sync.WaitGroup - ) - consumeWg.Add(2) - go func() { - defer consumeWg.Done() - for range concat.InitCh { - initEvents++ - } - }() - go func() { - defer consumeWg.Done() - for seg := range concat.SegCh { - segEvents++ - t.Logf("segment event %d: %d bytes", segEvents, len(seg)) - } - }() - - for i, p := range segPaths { - data, err := os.ReadFile(p) - require.NoError(t, err) - t.Logf("writing segment %d (%s, %d bytes)", i+1, filepath.Base(p), len(data)) - require.NoError(t, concat.Write(data)) - } - - require.NoError(t, concat.Close()) - consumeWg.Wait() - - require.Equal(t, 1, initEvents, "expected exactly one init event") - require.Equal(t, len(segPaths), segEvents, - "expected one segment event per input file (got %d for %d inputs)", - segEvents, len(segPaths)) -} diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go index ca8277b41..22473c716 100644 --- a/pkg/s3/s3.go +++ b/pkg/s3/s3.go @@ -33,17 +33,18 @@ type Recorder interface { } // S3Uploader manages streaming multipart uploads to an S3-compatible endpoint. -// Full fMP4 archives are fed via AddSegment. They are run through a muxl -// Concatenator to strip duplicate init segments, then uploaded as a -// multipart upload. Every cutoverEvery, the current upload is completed -// and a new one begins. +// Bare canonical MUXL segments are fed via AddSegment; they concatenate +// directly (the format is naively concatenable). A single init segment, +// synthesized from the first segment, is prepended to each multipart object so +// each object is a valid standalone MP4. Every cutoverEvery, the current upload +// is completed and a new one begins. type S3Uploader struct { client *s3.Client bucket string cutoverEvery time.Duration keyPrefix string // e.g. "did:plc:abc123/" userDID string - concat *muxl.Concatenator + segCh chan []byte // bare canonical MUXL segments awaiting upload done chan error recorder Recorder } @@ -84,14 +85,13 @@ func NewS3Uploader(cfg Config, userDID, keyPrefix string, cutoverEvery time.Dura if cutoverEvery == 0 { cutoverEvery = DefaultCutoverEvery } - concat := muxl.NewConcatenator(ctx) u := &S3Uploader{ client: client, bucket: cfg.Bucket, cutoverEvery: cutoverEvery, keyPrefix: keyPrefix, userDID: userDID, - concat: concat, + segCh: make(chan []byte, 16), done: make(chan error, 1), recorder: recorder, } @@ -99,25 +99,30 @@ func NewS3Uploader(cfg Config, userDID, keyPrefix string, cutoverEvery time.Dura return u } -// AddSegment feeds a full fMP4 archive (init+segments) to the concatenator -// for processing and upload. +// AddSegment feeds one bare canonical MUXL segment (uuid+moof+mdat per track) +// for upload. The bytes are copied, so the caller may reuse its buffer. func (u *S3Uploader) AddSegment(ctx context.Context, data []byte) error { - return u.concat.Write(data) + seg := append([]byte(nil), data...) + select { + case u.segCh <- seg: + return nil + case <-ctx.Done(): + return ctx.Err() + } } // Close signals that no more segments will be added, waits for all // in-flight uploads to complete, and returns any error. func (u *S3Uploader) Close(ctx context.Context) error { - closeErr := u.concat.Close() - uploadErr := <-u.done - if uploadErr != nil { + close(u.segCh) + if uploadErr := <-u.done; uploadErr != nil { return fmt.Errorf("error uploading: %w", uploadErr) } - return closeErr + return nil } -// uploadLoop reads init and segment events from the concatenator and manages -// multipart uploads. Runs until the concatenator's channels are closed. +// uploadLoop reads bare segments off segCh and manages multipart uploads. +// Runs until segCh is closed (Close) or ctx is canceled. func (u *S3Uploader) uploadLoop(ctx context.Context) { ctx = log.WithLogValues(ctx, "func", "s3.uploadLoop") var initSeg []byte @@ -257,18 +262,9 @@ func (u *S3Uploader) uploadLoop(ctx context.Context) { var err error for err == nil { select { - case init, ok := <-u.concat.InitCh: + case seg, ok := <-u.segCh: if !ok { - u.concat.InitCh = nil - continue - } - initSeg = init - log.Debug(ctx, "received init segment for S3 upload", "size", len(init)) - - case seg, ok := <-u.concat.SegCh: - log.Debug(ctx, "received segment for S3 upload", "size", len(seg)) - if !ok { - // Concatenator is done, complete any in-progress upload + // No more segments; complete any in-progress upload. err = completeUpload() if err != nil { err = fmt.Errorf("error completing upload: %w", err) @@ -276,6 +272,19 @@ func (u *S3Uploader) uploadLoop(ctx context.Context) { u.done <- err return } + log.Debug(ctx, "received segment for S3 upload", "size", len(seg)) + // Synthesize the init segment once, from the first segment's + // embedded catalog; it's prepended to each multipart object. + if initSeg == nil { + var initBuf bytes.Buffer + if werr := muxl.RunMuxlWrapInit(ctx, bytes.NewReader(seg), &initBuf); werr != nil { + err = fmt.Errorf("synthesizing init segment: %w", werr) + u.done <- err + return + } + initSeg = initBuf.Bytes() + log.Debug(ctx, "synthesized init segment for S3 upload", "size", len(initSeg)) + } if err = handleSegment(seg); err != nil { log.Error(ctx, "error handling segment", "error", err) }