From 68f6635706e131e2cde32fc7b6d4432dff4ba0f2 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 25 Jun 2026 11:18:46 -0700 Subject: [PATCH] live-vod: only record published segments to S3; cut over on unpublish S3 recording ran before the published gate, so segments arriving after a livestream ended (or before it started) kept appending to the same rolling object, tagged with the just-ended livestream URI, and that object wasn't completed until the cutoverEvery timer or stream teardown. That produced either "no recorded S3 segments" (finalize counts only completed objects) or an over-long VOD (the post-stop tail got absorbed into the recording). Gate s3Upload on notif.Metadata.Published (an un-ended place.stream.livestream record), and on unpublished segments call the new S3Uploader.Cutover to complete the current object immediately so it's finalize-able right away. Co-Authored-By: Claude Opus 4.8 (1M context) --- pkg/director/s3_upload.go | 14 ++++++++++ pkg/director/stream_session.go | 13 ++++++++- pkg/s3/s3.go | 48 +++++++++++++++++++++++++++++----- pkg/s3/uploader_test.go | 29 ++++++++++++++++++++ 4 files changed, 97 insertions(+), 7 deletions(-) diff --git a/pkg/director/s3_upload.go b/pkg/director/s3_upload.go index 711d4ef7..5cc73019 100644 --- a/pkg/director/s3_upload.go +++ b/pkg/director/s3_upload.go @@ -49,6 +49,20 @@ func (ss *StreamSession) s3Upload(ctx context.Context, notif *media.NewSegmentNo }) } +// s3Cutover completes the current live-recording object so it's immediately +// finalize-able. Called when a segment arrives that is not part of a live +// (published) stream — i.e. the livestream just ended, or hasn't started — so +// the recording is closed out promptly rather than lingering un-completed until +// stream teardown. +func (ss *StreamSession) s3Cutover(ctx context.Context) { + if ss.s3Uploader == nil { + return + } + ss.Go(ctx, func() error { + return ss.s3Uploader.Cutover(ctx) + }) +} + func (ss *StreamSession) s3Close(ctx context.Context) { if ss.s3Uploader == nil { return diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 5642b8af..8f39be58 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -242,7 +242,18 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment return fmt.Errorf("could not convert segment to streamplace segment: %w", err) } - ss.s3Upload(ctx, notif) + // Record to S3 for live-to-VOD only while the stream is live (an un-ended + // place.stream.livestream record exists -> the segment is published). + // Segments that arrive before "go live" or after the livestream ends are + // NOT recorded: pushing them produced over-long VODs (the post-stop tail + // got absorbed into the recording). The moment publishing stops we complete + // the current object so it's immediately finalize-able instead of lingering + // un-completed (which made finalize report "no recorded S3 segments"). + if notif.Metadata.Published { + ss.s3Upload(ctx, notif) + } else { + ss.s3Cutover(ctx) + } ss.bus.Publish(spseg.Creator, spseg) ss.Go(ctx, func() error { diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go index 008ee1cc..eb475578 100644 --- a/pkg/s3/s3.go +++ b/pkg/s3/s3.go @@ -57,7 +57,7 @@ type S3Uploader struct { cutoverEvery time.Duration keyPrefix string // e.g. "did:plc:abc123/" userDID string - segCh chan []byte // bare canonical MUXL segments awaiting upload + segCh chan uploadCmd // bare canonical MUXL segments / cutover requests done chan error recorder Recorder @@ -136,7 +136,7 @@ func newS3Uploader(client uploadAPI, bucket, userDID, keyPrefix string, cutoverE cutoverEvery: cutoverEvery, keyPrefix: keyPrefix, userDID: userDID, - segCh: make(chan []byte, 16), + segCh: make(chan uploadCmd, 16), done: make(chan error, 1), recorder: recorder, } @@ -144,6 +144,14 @@ func newS3Uploader(client uploadAPI, bucket, userDID, keyPrefix string, cutoverE return u } +// uploadCmd is one item on segCh: either a segment to append (seg != nil) or a +// request to complete the current object now (cutover). Both travel the same +// channel so a cutover stays FIFO-ordered behind the segments queued before it. +type uploadCmd struct { + seg []byte // bare canonical MUXL segment to append; nil for a cutover + cutover bool // complete the current object now (see Cutover) +} + // 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. It is // the caller's responsibility not to call AddSegment concurrently with Close @@ -156,7 +164,26 @@ func (u *S3Uploader) AddSegment(ctx context.Context, data []byte) error { } seg := append([]byte(nil), data...) select { - case u.segCh <- seg: + case u.segCh <- uploadCmd{seg: seg}: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +// Cutover completes the current in-progress object (if any) so it becomes a +// finalize-able, completed S3 segment, without tearing the uploader down — the +// next AddSegment simply starts a fresh object. It's used when a livestream +// ends (or the stream goes unpublished): the recording is closed out promptly +// instead of waiting for the cutoverEvery timer or stream teardown, so finalize +// can find the completed objects right away. A no-op if there's no current +// object, and a no-op after Close. +func (u *S3Uploader) Cutover(ctx context.Context) error { + if u.closed.Load() { + return nil + } + select { + case u.segCh <- uploadCmd{cutover: true}: return nil case <-ctx.Done(): return ctx.Err() @@ -325,7 +352,7 @@ func (u *S3Uploader) uploadLoop(ctx context.Context) { var err error for err == nil { select { - case seg, ok := <-u.segCh: + case cmd, ok := <-u.segCh: if !ok { // No more segments; complete any in-progress upload. err = completeUpload() @@ -335,8 +362,17 @@ func (u *S3Uploader) uploadLoop(ctx context.Context) { u.done <- err return } - log.Debug(ctx, "received segment for S3 upload", "size", len(seg)) - if err = handleSegment(seg); err != nil { + if cmd.cutover { + // Close out the current object so it's immediately finalize-able + // (e.g. the livestream just ended). No-op if nothing is in flight. + if err = completeUpload(); err != nil { + err = fmt.Errorf("error completing upload on cutover: %w", err) + log.Error(ctx, "error completing upload on cutover", "error", err) + } + continue + } + log.Debug(ctx, "received segment for S3 upload", "size", len(cmd.seg)) + if err = handleSegment(cmd.seg); err != nil { log.Error(ctx, "error handling segment", "error", err) } diff --git a/pkg/s3/uploader_test.go b/pkg/s3/uploader_test.go index 69a498e0..523fd3da 100644 --- a/pkg/s3/uploader_test.go +++ b/pkg/s3/uploader_test.go @@ -116,6 +116,35 @@ func TestS3UploaderCutoverOnLivestreamChange(t *testing.T) { require.NotEqual(t, keys[0], keys[1], "rolled-over objects must have distinct keys") } +// TestS3UploaderCutoverCompletesObject proves Cutover closes out the current +// object so the next segment starts a fresh one. This is what makes a recording +// finalize-able the moment a livestream ends, rather than lingering until the +// cutoverEvery timer or stream teardown. cutoverEvery is huge so only the +// explicit Cutover can trigger the rollover. +func TestS3UploaderCutoverCompletesObject(t *testing.T) { + fc := &fakeUploadAPI{} + rec := &fakeRecorder{} + u := newS3Uploader(fc, "bucket", "did:plc:test", "did:plc:test/", time.Hour, rec) + + ctx := context.Background() + seg := make([]byte, 1024) // under minPartSize: buffered until the object completes + + u.SetLivestreamURI("at://A") + require.NoError(t, u.AddSegment(ctx, seg)) + waitForStarts(t, rec, 1) // object 1 + + require.NoError(t, u.Cutover(ctx)) // completes object 1 + require.NoError(t, u.AddSegment(ctx, seg)) + waitForStarts(t, rec, 2) // object 2 + + require.NoError(t, u.Close(ctx)) + + require.Equal(t, 2, rec.count(), "Cutover must close the object so the next segment starts a new one") + keys := rec.startKeys() + require.Len(t, keys, 2) + require.NotEqual(t, keys[0], keys[1], "post-cutover object must have a distinct key") +} + // TestS3UploaderCloseIdempotent exercises the lifecycle fix that re-enabled // live S3 upload: Close must be safe to call repeatedly and concurrently (it // was a plain close(segCh) before, which panicked on the second call), and -- 2.51.2