diff --git a/js/docs/src/content/docs/lex-reference/media/place-stream-media-finalizelivestream.md b/js/docs/src/content/docs/lex-reference/media/place-stream-media-finalizelivestream.md new file mode 100644 index 00000000..79e8f4f0 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/media/place-stream-media-finalizelivestream.md @@ -0,0 +1,99 @@ +--- +title: place.stream.media.finalizeLivestream +description: Reference for the place.stream.media.finalizeLivestream lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Turn a finished livestream into a VOD. The server concatenates the MUXL segments it recorded for the livestream into a single content blob, derives the playback sidecars, and publishes the place.stream.media.track records — the same end state as a finished upload. Returns an uploadId the client polls with place.stream.media.getUploadStatus and then publishes with place.stream.media.publishVideo, exactly as for a resumable upload. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ------------ | -------- | ----- | ----------------------------------------------------------------------------------------------------------- | ---------------- | +| `livestream` | `string` | ✅ | AT-URI of the place.stream.livestream record to finalize into a VOD. Must belong to the authenticated user. | Format: `at-uri` | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ---------- | -------- | ----- | --------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------- | +| `uploadId` | `string` | ✅ | Identifier for the finalize job. Poll place.stream.media.getUploadStatus with it; once status is 'done', create the video with place.stream.media.publishVideo. | | + +**Possible Errors:** + +- `LivestreamNotFound`: No livestream with the given URI is known, or it does not belong to the authenticated user. +- `NoRecording`: The livestream has no recorded MUXL segments to finalize (recording was not enabled, or none completed). + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.media.finalizeLivestream", + "defs": { + "main": { + "type": "procedure", + "description": "Turn a finished livestream into a VOD. The server concatenates the MUXL segments it recorded for the livestream into a single content blob, derives the playback sidecars, and publishes the place.stream.media.track records — the same end state as a finished upload. Returns an uploadId the client polls with place.stream.media.getUploadStatus and then publishes with place.stream.media.publishVideo, exactly as for a resumable upload.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["livestream"], + "properties": { + "livestream": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the place.stream.livestream record to finalize into a VOD. Must belong to the authenticated user." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uploadId"], + "properties": { + "uploadId": { + "type": "string", + "description": "Identifier for the finalize job. Poll place.stream.media.getUploadStatus with it; once status is 'done', create the video with place.stream.media.publishVideo." + } + } + } + }, + "errors": [ + { + "name": "LivestreamNotFound", + "description": "No livestream with the given URI is known, or it does not belong to the authenticated user." + }, + { + "name": "NoRecording", + "description": "The livestream has no recorded MUXL segments to finalize (recording was not enabled, or none completed)." + } + ] + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index dd385ba8..7f440acb 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -2286,6 +2286,77 @@ } } }, + "/xrpc/place.stream.media.finalizeLivestream": { + "post": { + "summary": "Turn a finished livestream into a VOD. The server concatenates the MUXL segments it recorded for the livestream into a single content blob, derives the playback sidecars, and publishes the place.stream.media.track records — the same end state as a finished upload. Returns an uploadId the client polls with place.stream.media.getUploadStatus and then publishes with place.stream.media.publishVideo, exactly as for a resumable upload.", + "operationId": "place.stream.media.finalizeLivestream", + "tags": ["place.stream.media"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "uploadId": { + "type": "string", + "description": "Identifier for the finalize job. Poll place.stream.media.getUploadStatus with it; once status is 'done', create the video with place.stream.media.publishVideo." + } + }, + "required": ["uploadId"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "LivestreamNotFound" + }, + { + "const": "NoRecording" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "livestream": { + "type": "string", + "description": "AT-URI of the place.stream.livestream record to finalize into a VOD. Must belong to the authenticated user.", + "format": "uri" + } + }, + "required": ["livestream"] + } + } + } + } + } + }, "/xrpc/place.stream.media.getUploadStatus": { "get": { "summary": "Get the processing status of a previously created upload. Only accessible by the DID that created the upload.", diff --git a/lexicons/place/stream/media/finalizeLivestream.json b/lexicons/place/stream/media/finalizeLivestream.json new file mode 100644 index 00000000..5c8c96af --- /dev/null +++ b/lexicons/place/stream/media/finalizeLivestream.json @@ -0,0 +1,47 @@ +{ + "lexicon": 1, + "id": "place.stream.media.finalizeLivestream", + "defs": { + "main": { + "type": "procedure", + "description": "Turn a finished livestream into a VOD. The server concatenates the MUXL segments it recorded for the livestream into a single content blob, derives the playback sidecars, and publishes the place.stream.media.track records — the same end state as a finished upload. Returns an uploadId the client polls with place.stream.media.getUploadStatus and then publishes with place.stream.media.publishVideo, exactly as for a resumable upload.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["livestream"], + "properties": { + "livestream": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the place.stream.livestream record to finalize into a VOD. Must belong to the authenticated user." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uploadId"], + "properties": { + "uploadId": { + "type": "string", + "description": "Identifier for the finalize job. Poll place.stream.media.getUploadStatus with it; once status is 'done', create the video with place.stream.media.publishVideo." + } + } + } + }, + "errors": [ + { + "name": "LivestreamNotFound", + "description": "No livestream with the given URI is known, or it does not belong to the authenticated user." + }, + { + "name": "NoRecording", + "description": "The livestream has no recorded MUXL segments to finalize (recording was not enabled, or none completed)." + } + ] + } + } +} diff --git a/pkg/blob/s3.go b/pkg/blob/s3.go index 22f0bd27..73dc3bf1 100644 --- a/pkg/blob/s3.go +++ b/pkg/blob/s3.go @@ -40,6 +40,12 @@ func (s *S3Store) URL(key string) string { return "s3://" + s.bucket + "/" + key func (s *S3Store) Bucket() string { return s.bucket } +// Client exposes the underlying S3 client so callers needing operations beyond +// the blob.Store interface (e.g. the server-side multi-object concat in +// pkg/s3.ConcatWithHeader for live-to-VOD finalize) can reuse the same +// configured client + bucket. +func (s *S3Store) Client() *awss3.Client { return s.client } + func (s *S3Store) Open(ctx context.Context, key string) (Reader, error) { ra, err := s3pkg.NewReaderAt(ctx, s.client, s.bucket, key) if err != nil { diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 31f9625f..cebfcb90 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -413,6 +413,21 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu Location: t.Location, }) }) + // Live-to-VOD finalize. Resolves the streamer's live signing key here + // (where the model is in scope) so pkg/vod stays free of pkg/model, then + // concatenates the recorded MUXL objects into a VOD. + state.SetLivestreamVODFinalizer(func(ctx context.Context, t statedb.FinalizeLivestreamVODTask) (string, error) { + signingKey, err := resolveLiveSigningKey(mod, t.RepoDID) + if err != nil { + return "", fmt.Errorf("finalize-livestream-vod: resolve signing key: %w", err) + } + return vod.FinalizeLivestreamVOD(ctx, cli, state, vodStore, vod.FinalizeInput{ + UploadID: t.UploadID, + RepoDID: t.RepoDID, + LivestreamURI: t.LivestreamURI, + SigningKey: signingKey, + }) + }) // View-count aggregator runs the log → record pipeline for one // window. Same function-pointer pattern as the VOD processor so // statedb stays free of viewlog's transitive deps. The scheduler @@ -1139,3 +1154,29 @@ func makeMigrateCommand(build *config.BuildFlags) *urfavecli.Command { }, } } + +// resolveLiveSigningKey returns the did:key whose private half signed a +// streamer's live segments, for stamping on live-to-VOD place.stream.media.track +// records so playback can verify them. It picks the most recently created +// non-revoked place.stream.key for the repo; streamers normally have exactly +// one. Errors if the repo has no active signing key. +func resolveLiveSigningKey(mod model.Model, repoDID string) (string, error) { + keys, err := mod.GetSigningKeysForRepo(repoDID) + if err != nil { + return "", err + } + var best *model.SigningKey + for i := range keys { + k := &keys[i] + if k.RevokedAt != nil { + continue + } + if best == nil || k.CreatedAt.After(best.CreatedAt) { + best = k + } + } + if best == nil { + return "", fmt.Errorf("no active signing key for repo %s", repoDID) + } + return best.DID, nil +} diff --git a/pkg/director/s3_upload.go b/pkg/director/s3_upload.go index a2dad249..711d4ef7 100644 --- a/pkg/director/s3_upload.go +++ b/pkg/director/s3_upload.go @@ -2,13 +2,16 @@ package director import ( "context" - "time" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/s3" ) +// liveRecPrefix namespaces in-progress livestream recordings in the S3 bucket +// (keys are liveRecPrefix + / + .m4s). +const liveRecPrefix = "live-rec/" + func (ss *StreamSession) maybeStartS3Upload(ctx context.Context, repoDID string) { if !ss.cli.S3Configured() { return @@ -20,8 +23,18 @@ func (ss *StreamSession) maybeStartS3Upload(ctx context.Context, repoDID string) SecretAccessKey: ss.cli.S3SecretAccessKey, Region: ss.cli.S3Region, } - keyPrefix := repoDID + "/" - ss.s3Uploader = s3.NewS3Uploader(cfg, repoDID, keyPrefix, time.Minute, ss.statefulDB) + // live-rec/ namespaces the in-progress livestream recordings away from the + // finalized VOD blobs (blobs/) and anything else in the bucket. + keyPrefix := liveRecPrefix + repoDID + "/" + ss.s3Uploader = s3.NewS3Uploader(cfg, repoDID, keyPrefix, s3.DefaultCutoverEvery, ss.statefulDB) + // Best-effort initial resolve of the livestream URI so the very first + // object is tagged. The director treats "latest livestream for repo" as + // the current stream everywhere (notification blast, idle finalize), so we + // do the same here; NewSegment refreshes it once the stream's own record is + // indexed, in case a prior stream was momentarily still "latest". + if ls, err := ss.mod.GetLatestLivestreamForRepo(repoDID); err == nil && ls != nil { + ss.s3Uploader.SetLivestreamURI(ls.URI) + } log.Log(ctx, "S3 upload enabled", "bucket", ss.cli.S3Bucket, "endpoint", ss.cli.S3Endpoint) } diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 3a8aba0e..5642b8af 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -126,7 +126,7 @@ func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotif allRenditions = append([]renditions.Rendition{sourceRendition}, allRenditions...) allRenditions = append(allRenditions, renditions.AudioRendition) - // ss.maybeStartS3Upload(ctx, notif.Segment.RepoDID) + ss.maybeStartS3Upload(ctx, notif.Segment.RepoDID) close(ss.started) @@ -155,9 +155,13 @@ func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotif case <-ss.segmentChan: // reset timer case <-ctx.Done(): - // Signal all background workers to stop + // Drain all in-flight session goroutines (including the per-segment + // AddSegment senders) BEFORE closing the uploader, so no segment + // send can race the uploader's channel close. s3Close then flushes + // the final object on the uploader's own (still-live) context. + err := ss.g.Wait() ss.s3Close(ctx) - return ss.g.Wait() + return err // case <-time.After(time.Minute * 1): case <-time.After(ss.cli.StreamSessionTimeout): log.Log(ctx, "stream session timeout, shutting down", "timeout", ss.cli.StreamSessionTimeout) @@ -238,7 +242,7 @@ 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) + ss.s3Upload(ctx, notif) ss.bus.Publish(spseg.Creator, spseg) ss.Go(ctx, func() error { @@ -299,6 +303,12 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment log.Warn(ctx, "no livestream found, skipping notification blast", "repoDID", spseg.Creator) return nil } + // Refresh the S3 uploader's livestream tag now that this stream's + // own record is indexed (it may not have been when the uploader + // started), so all objects are attributed to the right stream. + if ss.s3Uploader != nil { + ss.s3Uploader.SetLivestreamURI(livestreamModel.URI) + } lsv, err := livestreamModel.ToLivestreamView() if err != nil { return fmt.Errorf("failed to convert livestream to streamplace livestream: %w", err) diff --git a/pkg/s3/concat.go b/pkg/s3/concat.go new file mode 100644 index 00000000..d4f8b422 --- /dev/null +++ b/pkg/s3/concat.go @@ -0,0 +1,250 @@ +package s3 + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/trace" + "stream.place/streamplace/pkg/log" +) + +// ErrConcatPartTooSmall is returned by ConcatWithHeader when a non-final source +// object is smaller than S3's 5 MB minimum part size and therefore can't be a +// server-side UploadPartCopy part. The caller is expected to fall back to a +// download-and-rewrite assembly. In production (10-minute cutover objects) this +// never triggers; it only guards against pathologically short streams. +var ErrConcatPartTooSmall = errors.New("s3 concat: non-final source object below 5MB minimum part size") + +// concatAPI is the subset of *s3.Client that ConcatWithHeader uses. Pulled out +// so tests can inject a fake; *s3.Client satisfies it. +type concatAPI interface { + copyAPI + GetObject(context.Context, *s3.GetObjectInput, ...func(*s3.Options)) (*s3.GetObjectOutput, error) + UploadPart(context.Context, *s3.UploadPartInput, ...func(*s3.Options)) (*s3.UploadPartOutput, error) +} + +// ConcatWithHeader assembles destKey = header ++ srcKeys[0] ++ srcKeys[1] ++ … +// (the source objects in order) as a single multipart object, transferring +// almost everything server-side. +// +// The layout exists to give live-to-VOD finalize a content blob shaped exactly +// like an uploaded VOD ([init][segments…]) without re-uploading the (possibly +// many-GB) stream: the synthesized init header rides in part 1, and the bulk of +// the bytes are copied with UploadPartCopy and never pass through this process. +// +// Part 1 (uploaded) is header plus just enough leading source bytes to clear +// S3's 5 MB minimum, so it costs ~5 MB of memory + one upload. Every remaining +// source object becomes one server-side UploadPartCopy part. S3 requires every +// part except the last to be ≥5 MB; a non-final source object below that bound +// can't be copied as its own part, so ConcatWithHeader returns +// ErrConcatPartTooSmall and the caller falls back to a full rewrite. +func ConcatWithHeader(ctx context.Context, client *s3.Client, bucket string, header []byte, srcKeys []string, dstKey, contentType string) error { + return concatWithHeader(ctx, client, bucket, header, srcKeys, dstKey, contentType) +} + +func concatWithHeader(ctx context.Context, client concatAPI, bucket string, header []byte, srcKeys []string, dstKey, contentType string) error { + ctx = log.WithLogValues(ctx, "func", "s3.ConcatWithHeader") + ctx, span := s3Tracer.Start(ctx, "s3.ConcatWithHeader", trace.WithAttributes( + attribute.String("bucket", bucket), + attribute.String("dst_key", dstKey), + attribute.Int("src_count", len(srcKeys)), + attribute.Int("header_bytes", len(header)), + )) + defer span.End() + + if len(srcKeys) == 0 { + return fmt.Errorf("s3 concat: no source objects") + } + + // Object sizes drive both the part-1 fill and the per-object copy parts. + sizes := make([]int64, len(srcKeys)) + for i, key := range srcKeys { + head, err := client.HeadObject(ctx, &s3.HeadObjectInput{ + Bucket: aws.String(bucket), + Key: aws.String(key), + }) + if err != nil { + span.RecordError(err) + return fmt.Errorf("head s3://%s/%s: %w", bucket, key, err) + } + sizes[i] = aws.ToInt64(head.ContentLength) + } + + create := &s3.CreateMultipartUploadInput{ + Bucket: aws.String(bucket), + Key: aws.String(dstKey), + } + if contentType != "" { + create.ContentType = aws.String(contentType) + } + resp, err := client.CreateMultipartUpload(ctx, create) + if err != nil { + span.RecordError(err) + return fmt.Errorf("create multipart upload s3://%s/%s: %w", bucket, dstKey, err) + } + uploadID := aws.ToString(resp.UploadId) + + abort := func() { + // Use ctx (not a cancelled child) so cleanup runs; best-effort. + _, _ = client.AbortMultipartUpload(ctx, &s3.AbortMultipartUploadInput{ + Bucket: aws.String(bucket), + Key: aws.String(dstKey), + UploadId: aws.String(uploadID), + }) + } + + var parts []types.CompletedPart + var partNum int32 + + // --- Part 1: header + leading source bytes until >= minPartSize. --- + // We pull only enough source bytes to clear the 5 MB floor (bounding memory + // to ~minPartSize), tracking which object we stopped in so the copy phase + // resumes from exactly there. + buf := bytes.NewBuffer(make([]byte, 0, 2*minPartSize+len(header))) + buf.Write(header) + idx := 0 // index of the source object the copy phase resumes at + var consumedInIdx int64 // bytes of srcKeys[idx] already pulled into part 1 + for idx < len(srcKeys) && int64(buf.Len()) < minPartSize { + remaining := int64(minPartSize) - int64(buf.Len()) + avail := sizes[idx] - consumedInIdx + take := remaining + if take > avail { + take = avail + } + // Look-ahead: a partial take that leaves this (non-final) object with a + // sub-5MB remainder would make that remainder an illegal small middle + // copy part. Absorb the whole object into part 1 instead — costing a + // little extra upload but keeping every copy part legal. + if take < avail && (avail-take) < int64(minPartSize) && idx != len(srcKeys)-1 { + take = avail + } + if take > 0 { + body, err := getRange(ctx, client, bucket, srcKeys[idx], consumedInIdx, consumedInIdx+take-1) + if err != nil { + abort() + span.RecordError(err) + return err + } + buf.Write(body) + consumedInIdx += take + } + if consumedInIdx >= sizes[idx] { + idx++ + consumedInIdx = 0 + } + } + + partNum++ + if err := uploadPart(ctx, client, bucket, dstKey, uploadID, partNum, buf.Bytes(), &parts); err != nil { + abort() + span.RecordError(err) + span.SetStatus(codes.Error, "upload_part_1") + return err + } + + // --- Remaining objects: one server-side UploadPartCopy part each. --- + for i := idx; i < len(srcKeys); i++ { + start := int64(0) + if i == idx { + start = consumedInIdx // resume mid-object if part 1 stopped here + } + if start >= sizes[i] { + continue // wholly absorbed into part 1 + } + partSize := sizes[i] - start + isLast := i == len(srcKeys)-1 + if !isLast && partSize < minPartSize { + abort() + err := fmt.Errorf("%w (object %s contributes %d bytes)", ErrConcatPartTooSmall, srcKeys[i], partSize) + span.RecordError(err) + return err + } + if partSize > maxCopyObjectSize { + abort() + err := fmt.Errorf("s3 concat: object %s range %d bytes exceeds single-part copy limit", srcKeys[i], partSize) + span.RecordError(err) + return err + } + partNum++ + res, err := client.UploadPartCopy(ctx, &s3.UploadPartCopyInput{ + Bucket: aws.String(bucket), + Key: aws.String(dstKey), + UploadId: aws.String(uploadID), + PartNumber: aws.Int32(partNum), + CopySource: aws.String(bucket + "/" + srcKeys[i]), + CopySourceRange: aws.String(fmt.Sprintf("bytes=%d-%d", start, sizes[i]-1)), + }) + if err != nil { + abort() + span.RecordError(err) + span.SetStatus(codes.Error, "upload_part_copy") + return fmt.Errorf("upload part copy %d s3://%s/%s: %w", partNum, bucket, srcKeys[i], err) + } + if res.CopyPartResult == nil { + abort() + return fmt.Errorf("upload part copy %d s3://%s/%s: missing CopyPartResult", partNum, bucket, srcKeys[i]) + } + parts = append(parts, types.CompletedPart{ + ETag: res.CopyPartResult.ETag, + PartNumber: aws.Int32(partNum), + }) + } + + if _, err := client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{ + Bucket: aws.String(bucket), + Key: aws.String(dstKey), + UploadId: aws.String(uploadID), + MultipartUpload: &types.CompletedMultipartUpload{Parts: parts}, + }); err != nil { + abort() + span.RecordError(err) + span.SetStatus(codes.Error, "complete") + return fmt.Errorf("complete multipart upload s3://%s/%s: %w", bucket, dstKey, err) + } + log.Log(ctx, "completed S3 header concat", "bucket", bucket, "key", dstKey, "parts", len(parts)) + return nil +} + +// getRange reads [start,end] (inclusive) of an object into memory. +func getRange(ctx context.Context, client concatAPI, bucket, key string, start, end int64) ([]byte, error) { + out, err := client.GetObject(ctx, &s3.GetObjectInput{ + Bucket: aws.String(bucket), + Key: aws.String(key), + Range: aws.String(fmt.Sprintf("bytes=%d-%d", start, end)), + }) + if err != nil { + return nil, fmt.Errorf("get s3://%s/%s bytes=%d-%d: %w", bucket, key, start, end, err) + } + defer out.Body.Close() + body, err := io.ReadAll(out.Body) + if err != nil { + return nil, fmt.Errorf("read s3://%s/%s body: %w", bucket, key, err) + } + return body, nil +} + +func uploadPart(ctx context.Context, client concatAPI, bucket, dstKey, uploadID string, partNum int32, body []byte, parts *[]types.CompletedPart) error { + res, err := client.UploadPart(ctx, &s3.UploadPartInput{ + Bucket: aws.String(bucket), + Key: aws.String(dstKey), + UploadId: aws.String(uploadID), + PartNumber: aws.Int32(partNum), + Body: bytes.NewReader(body), + }) + if err != nil { + return fmt.Errorf("upload part %d s3://%s/%s: %w", partNum, bucket, dstKey, err) + } + *parts = append(*parts, types.CompletedPart{ + ETag: res.ETag, + PartNumber: aws.Int32(partNum), + }) + return nil +} diff --git a/pkg/s3/concat_test.go b/pkg/s3/concat_test.go new file mode 100644 index 00000000..d02af9e2 --- /dev/null +++ b/pkg/s3/concat_test.go @@ -0,0 +1,214 @@ +package s3 + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "sort" + "strconv" + "strings" + "testing" + + "github.com/aws/aws-sdk-go-v2/aws" + awss3 "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" +) + +// fakeConcatS3 is an in-memory stand-in for the subset of S3 ConcatWithHeader +// uses. CompleteMultipartUpload reconstructs the destination object from the +// recorded parts in part-number order and enforces S3's rule that every part +// except the last is ≥5MB, so a ConcatWithHeader bug that emits a small middle +// part fails loudly here rather than only against real S3. +type fakeConcatS3 struct { + objects map[string][]byte + parts map[int32][]byte // partNumber -> bytes, for the single in-flight upload + assembled map[string][]byte + aborted bool +} + +func newFakeConcatS3(objects map[string][]byte) *fakeConcatS3 { + return &fakeConcatS3{ + objects: objects, + parts: map[int32][]byte{}, + assembled: map[string][]byte{}, + } +} + +func (f *fakeConcatS3) HeadObject(_ context.Context, in *awss3.HeadObjectInput, _ ...func(*awss3.Options)) (*awss3.HeadObjectOutput, error) { + b, ok := f.objects[aws.ToString(in.Key)] + if !ok { + return nil, fmt.Errorf("no such object %s", aws.ToString(in.Key)) + } + return &awss3.HeadObjectOutput{ContentLength: aws.Int64(int64(len(b)))}, nil +} + +func (f *fakeConcatS3) GetObject(_ context.Context, in *awss3.GetObjectInput, _ ...func(*awss3.Options)) (*awss3.GetObjectOutput, error) { + b, ok := f.objects[aws.ToString(in.Key)] + if !ok { + return nil, fmt.Errorf("no such object %s", aws.ToString(in.Key)) + } + start, end := parseRange(aws.ToString(in.Range), len(b)) + chunk := b[start : end+1] + return &awss3.GetObjectOutput{Body: io.NopCloser(bytes.NewReader(chunk))}, nil +} + +func (f *fakeConcatS3) CreateMultipartUpload(_ context.Context, _ *awss3.CreateMultipartUploadInput, _ ...func(*awss3.Options)) (*awss3.CreateMultipartUploadOutput, error) { + return &awss3.CreateMultipartUploadOutput{UploadId: aws.String("upload-1")}, nil +} + +func (f *fakeConcatS3) UploadPart(_ context.Context, in *awss3.UploadPartInput, _ ...func(*awss3.Options)) (*awss3.UploadPartOutput, error) { + body, _ := io.ReadAll(in.Body) + f.parts[aws.ToInt32(in.PartNumber)] = body + return &awss3.UploadPartOutput{ETag: aws.String(fmt.Sprintf("etag-%d", aws.ToInt32(in.PartNumber)))}, nil +} + +func (f *fakeConcatS3) UploadPartCopy(_ context.Context, in *awss3.UploadPartCopyInput, _ ...func(*awss3.Options)) (*awss3.UploadPartCopyOutput, error) { + // CopySource is "bucket/key"; strip the bucket prefix. + src := aws.ToString(in.CopySource) + key := src[strings.Index(src, "/")+1:] + b, ok := f.objects[key] + if !ok { + return nil, fmt.Errorf("copy from missing object %s", key) + } + start, end := parseRange(aws.ToString(in.CopySourceRange), len(b)) + f.parts[aws.ToInt32(in.PartNumber)] = append([]byte(nil), b[start:end+1]...) + return &awss3.UploadPartCopyOutput{ + CopyPartResult: &types.CopyPartResult{ETag: aws.String(fmt.Sprintf("etag-%d", aws.ToInt32(in.PartNumber)))}, + }, nil +} + +func (f *fakeConcatS3) CompleteMultipartUpload(_ context.Context, in *awss3.CompleteMultipartUploadInput, _ ...func(*awss3.Options)) (*awss3.CompleteMultipartUploadOutput, error) { + nums := make([]int32, 0, len(in.MultipartUpload.Parts)) + for _, p := range in.MultipartUpload.Parts { + nums = append(nums, aws.ToInt32(p.PartNumber)) + } + sort.Slice(nums, func(i, j int) bool { return nums[i] < nums[j] }) + var out []byte + for i, n := range nums { + body := f.parts[n] + isLast := i == len(nums)-1 + if !isLast && len(body) < minPartSize { + return nil, fmt.Errorf("EntityTooSmall: part %d is %d bytes (<5MB) and not last", n, len(body)) + } + out = append(out, body...) + } + f.assembled[aws.ToString(in.Key)] = out + return &awss3.CompleteMultipartUploadOutput{}, nil +} + +func (f *fakeConcatS3) AbortMultipartUpload(_ context.Context, _ *awss3.AbortMultipartUploadInput, _ ...func(*awss3.Options)) (*awss3.AbortMultipartUploadOutput, error) { + f.aborted = true + return &awss3.AbortMultipartUploadOutput{}, nil +} + +// CopyObject is required to satisfy copyAPI (embedded in concatAPI) but unused. +func (f *fakeConcatS3) CopyObject(_ context.Context, _ *awss3.CopyObjectInput, _ ...func(*awss3.Options)) (*awss3.CopyObjectOutput, error) { + return nil, errors.New("unexpected CopyObject") +} + +func parseRange(r string, total int) (int64, int64) { + // "bytes=start-end" + spec := strings.TrimPrefix(r, "bytes=") + parts := strings.SplitN(spec, "-", 2) + start, _ := strconv.ParseInt(parts[0], 10, 64) + end, _ := strconv.ParseInt(parts[1], 10, 64) + if end > int64(total)-1 { + end = int64(total) - 1 + } + return start, end +} + +func filled(b byte, n int) []byte { + out := make([]byte, n) + for i := range out { + out[i] = b + } + return out +} + +const mb = 1024 * 1024 + +func TestConcatWithHeader(t *testing.T) { + header := filled('H', 1024) + + cases := []struct { + name string + objects map[string][]byte + order []string + wantErr error + }{ + { + name: "large objects partial first part", + objects: map[string][]byte{ + "a": filled('a', 10*mb), + "b": filled('b', 8*mb), + "c": filled('c', 7*mb), + }, + order: []string{"a", "b", "c"}, + }, + { + name: "medium objects absorb whole first object", + objects: map[string][]byte{ + "a": filled('a', 6*mb), + "b": filled('b', 6*mb), + "c": filled('c', 6*mb), + }, + order: []string{"a", "b", "c"}, + }, + { + name: "single small object is the only part", + objects: map[string][]byte{ + "a": filled('a', 2*mb), + }, + order: []string{"a"}, + }, + { + name: "small last object is allowed", + objects: map[string][]byte{ + "a": filled('a', 8*mb), + "b": filled('b', 1*mb), + }, + order: []string{"a", "b"}, + }, + { + name: "small middle object falls back", + objects: map[string][]byte{ + "a": filled('a', 6*mb), + "b": filled('b', 2*mb), + "c": filled('c', 6*mb), + }, + order: []string{"a", "b", "c"}, + wantErr: ErrConcatPartTooSmall, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + fake := newFakeConcatS3(tc.objects) + err := concatWithHeader(context.Background(), fake, "bucket", header, tc.order, "dst", "video/mp4") + if tc.wantErr != nil { + if !errors.Is(err, tc.wantErr) { + t.Fatalf("want error %v, got %v", tc.wantErr, err) + } + if !fake.aborted { + t.Fatalf("expected multipart upload to be aborted on too-small part") + } + return + } + if err != nil { + t.Fatalf("concatWithHeader: %v", err) + } + // Expected = header ++ objects in order. + want := append([]byte(nil), header...) + for _, k := range tc.order { + want = append(want, tc.objects[k]...) + } + got := fake.assembled["dst"] + if !bytes.Equal(got, want) { + t.Fatalf("assembled mismatch: got %d bytes, want %d bytes", len(got), len(want)) + } + }) + } +} diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go index 22473c71..7d4c0f17 100644 --- a/pkg/s3/s3.go +++ b/pkg/s3/s3.go @@ -4,6 +4,8 @@ import ( "bytes" "context" "fmt" + "sync" + "sync/atomic" "time" "github.com/aws/aws-sdk-go-v2/aws" @@ -11,7 +13,6 @@ import ( "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/aws/aws-sdk-go-v2/service/s3/types" "stream.place/streamplace/pkg/log" - "stream.place/streamplace/pkg/muxl" ) // Config holds the configuration for an S3-compatible upload target. @@ -26,18 +27,21 @@ type Config struct { // Recorder is an optional persistence hook for S3Uploader. RecordStart is // called when a new multipart upload begins; the returned id is passed back // to RecordComplete when the upload is finalized. Implementations should -// tolerate nil contexts being passed. +// tolerate nil contexts being passed. livestreamURI ties the object to the +// livestream it belongs to (may be empty if not yet known) so the live-to-VOD +// finalize can later enumerate exactly the objects for one stream. type Recorder interface { - RecordStart(ctx context.Context, userDID, bucket, key string, started time.Time) (id string, err error) + RecordStart(ctx context.Context, userDID, bucket, key, livestreamURI string, started time.Time) (id string, err error) RecordComplete(ctx context.Context, id string, parts int32, size int64) error } // S3Uploader manages streaming multipart uploads to an S3-compatible endpoint. -// 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. +// Bare canonical MUXL segments are fed via AddSegment and written verbatim — +// no per-object init header — so the objects of one stream concatenate +// directly into a single MUXL byte stream. The live-to-VOD finalize prepends +// one synthesized init to the whole concatenation; keeping the objects bare is +// what lets it "concat fearlessly". Every cutoverEvery, the current upload is +// completed and a new one begins. type S3Uploader struct { client *s3.Client bucket string @@ -47,6 +51,28 @@ type S3Uploader struct { segCh chan []byte // bare canonical MUXL segments awaiting upload done chan error recorder Recorder + + mu sync.Mutex + livestreamURI string // guarded by mu; stamped on each S3Segment row + + closeOnce sync.Once + closeErr error + closed atomic.Bool +} + +// SetLivestreamURI records the livestream this stream's objects belong to. It +// may be called once the URI is resolved (it can be unknown when the uploader +// starts); subsequent multipart objects are stamped with it via RecordStart. +func (u *S3Uploader) SetLivestreamURI(uri string) { + u.mu.Lock() + u.livestreamURI = uri + u.mu.Unlock() +} + +func (u *S3Uploader) getLivestreamURI() string { + u.mu.Lock() + defer u.mu.Unlock() + return u.livestreamURI } // S3 requires each part except the last to be at least 5MB. @@ -100,8 +126,15 @@ func NewS3Uploader(cfg Config, userDID, keyPrefix string, cutoverEvery time.Dura } // 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. +// 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 +// (the StreamSession guarantees this by draining its goroutines before closing, +// see director.Start); the closed guard here is a best-effort backstop that +// turns a late call into an error rather than a send-on-closed-channel panic. func (u *S3Uploader) AddSegment(ctx context.Context, data []byte) error { + if u.closed.Load() { + return fmt.Errorf("s3 uploader closed") + } seg := append([]byte(nil), data...) select { case u.segCh <- seg: @@ -111,32 +144,39 @@ func (u *S3Uploader) AddSegment(ctx context.Context, data []byte) error { } } -// Close signals that no more segments will be added, waits for all -// in-flight uploads to complete, and returns any error. +// Close signals that no more segments will be added, waits for all in-flight +// uploads to complete, and returns any error. It is idempotent: repeated calls +// return the same result without re-closing the channel. The supplied ctx is +// unused for the wait (uploadLoop runs on its own context so it can flush the +// final object even after the session context is cancelled) but kept for API +// symmetry. func (u *S3Uploader) Close(ctx context.Context) error { - close(u.segCh) - if uploadErr := <-u.done; uploadErr != nil { - return fmt.Errorf("error uploading: %w", uploadErr) - } - return nil + u.closeOnce.Do(func() { + u.closed.Store(true) + close(u.segCh) + if uploadErr := <-u.done; uploadErr != nil { + u.closeErr = fmt.Errorf("error uploading: %w", uploadErr) + } + }) + return u.closeErr } // uploadLoop reads bare segments off segCh and manages multipart uploads. -// Runs until segCh is closed (Close) or ctx is canceled. +// Runs until segCh is closed (Close) or ctx is canceled. ctx is intentionally +// independent of the session context so a final object can still be completed +// after the stream tears down. func (u *S3Uploader) uploadLoop(ctx context.Context) { ctx = log.WithLogValues(ctx, "func", "s3.uploadLoop") - var initSeg []byte var current *activeUpload - // Helper: prepend init to buffer when starting a new upload startUpload := func() error { now := time.Now() - key := fmt.Sprintf("%s%s.mp4", u.keyPrefix, now.UTC().Format("2006-01-02T15-04-05")) + key := fmt.Sprintf("%s%s.m4s", u.keyPrefix, now.UTC().Format("2006-01-02T15-04-05")) resp, err := u.client.CreateMultipartUpload(ctx, &s3.CreateMultipartUploadInput{ Bucket: aws.String(u.bucket), Key: aws.String(key), - ContentType: aws.String("video/mp4"), + ContentType: aws.String("video/iso.segment"), }) if err != nil { return fmt.Errorf("creating multipart upload for %s: %w", key, err) @@ -148,16 +188,12 @@ func (u *S3Uploader) uploadLoop(ctx context.Context) { started: now, } if u.recorder != nil { - id, recErr := u.recorder.RecordStart(ctx, u.userDID, u.bucket, key, now) + id, recErr := u.recorder.RecordStart(ctx, u.userDID, u.bucket, key, u.getLivestreamURI(), now) if recErr != nil { log.Error(ctx, "recording S3 upload start", "key", key, "error", recErr) } current.recordID = id } - // Prepend init segment to the buffer so the file starts valid - if initSeg != nil { - current.buf = append(current.buf, initSeg...) - } log.Log(ctx, "started S3 multipart upload", "key", key) return nil } @@ -273,18 +309,6 @@ func (u *S3Uploader) uploadLoop(ctx context.Context) { 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) } diff --git a/pkg/s3/uploader_test.go b/pkg/s3/uploader_test.go new file mode 100644 index 00000000..8324cd63 --- /dev/null +++ b/pkg/s3/uploader_test.go @@ -0,0 +1,46 @@ +package s3 + +import ( + "context" + "sync" + "testing" + "time" +) + +// 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 +// AddSegment after Close must return an error rather than panic with +// "send on closed channel". No segments are added, so the upload loop completes +// without making any S3 calls — this stays a pure unit test. +func TestS3UploaderCloseIdempotent(t *testing.T) { + u := NewS3Uploader(Config{ + Region: "us-east-1", + Endpoint: "http://127.0.0.1:0", + Bucket: "test", + AccessKeyID: "k", + SecretAccessKey: "s", + }, "did:plc:test", "did:plc:test/", time.Minute, nil) + + const n = 4 + var wg sync.WaitGroup + errs := make([]error, n) + for i := 0; i < n; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + errs[i] = u.Close(context.Background()) + }(i) + } + wg.Wait() + for i, err := range errs { + if err != nil { + t.Fatalf("Close call %d returned error: %v", i, err) + } + } + + // A late AddSegment must be rejected, not panic on a closed channel. + if err := u.AddSegment(context.Background(), []byte("late")); err == nil { + t.Fatalf("AddSegment after Close should return an error") + } +} diff --git a/pkg/spxrpc/place_stream_media_finalizelivestream.go b/pkg/spxrpc/place_stream_media_finalizelivestream.go new file mode 100644 index 00000000..1e2c9feb --- /dev/null +++ b/pkg/spxrpc/place_stream_media_finalizelivestream.go @@ -0,0 +1,71 @@ +package spxrpc + +import ( + "context" + "net/http" + + "github.com/google/uuid" + "github.com/labstack/echo/v4" + "github.com/streamplace/oatproxy/pkg/oatproxy" + "stream.place/streamplace/pkg/statedb" + placestream "stream.place/streamplace/pkg/streamplace" +) + +// handlePlaceStreamMediaFinalizeLivestream turns a finished livestream into a +// VOD. It creates a synthetic Upload row and enqueues a background finalize +// task that concatenates the livestream's recorded MUXL objects into a content +// blob and publishes the track records — landing in the same place a finished +// resumable upload does, so the client polls getUploadStatus and then +// publishVideo with the returned uploadId, unchanged. +func (s *Server) handlePlaceStreamMediaFinalizeLivestream(ctx context.Context, body *placestream.MediaFinalizeLivestream_Input) (*placestream.MediaFinalizeLivestream_Output, error) { + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required") + } + if body.Livestream == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "livestream is required") + } + // A banned account can't mint new VODs (the playback gates would hide them + // regardless, but skip the work and the dead records). + if banned, err := s.accountBanned(session.DID); err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } else if banned { + return nil, echo.NewHTTPError(http.StatusForbidden, "account is not permitted to publish videos") + } + + ls, err := s.model.GetLivestream(body.Livestream) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + if ls == nil || ls.RepoDID != session.DID { + // Don't distinguish "not found" from "not yours" — same response. + return nil, echo.NewHTTPError(http.StatusNotFound, "livestream not found") + } + + // Synthetic Upload row so the client reuses the getUploadStatus / + // publishVideo flow it already has for resumable uploads. + uu, err := uuid.NewV7() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + uploadID := uu.String() + if err := s.statefulDB.CreateUpload(ctx, &statedb.Upload{ + ID: uploadID, + RepoDID: session.DID, + MimeType: "video/mp4", + Backend: "live", + Location: body.Livestream, + }); err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + + if _, err := s.statefulDB.EnqueueTask(ctx, statedb.TaskFinalizeLivestreamVOD, statedb.FinalizeLivestreamVODTask{ + UploadID: uploadID, + RepoDID: session.DID, + LivestreamURI: body.Livestream, + }, statedb.WithTaskKey("finalize-vod:"+uploadID)); err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + + return &placestream.MediaFinalizeLivestream_Output{UploadId: uploadID}, nil +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index abbddb30..77c15f0f 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -305,6 +305,7 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error { e.POST("/xrpc/place.stream.live.startLivestream", s.HandlePlaceStreamLiveStartLivestream) e.POST("/xrpc/place.stream.live.stopLivestream", s.HandlePlaceStreamLiveStopLivestream) e.POST("/xrpc/place.stream.media.createUpload", s.HandlePlaceStreamMediaCreateUpload) + e.POST("/xrpc/place.stream.media.finalizeLivestream", s.HandlePlaceStreamMediaFinalizeLivestream) e.GET("/xrpc/place.stream.media.getUploadStatus", s.HandlePlaceStreamMediaGetUploadStatus) e.GET("/xrpc/place.stream.media.getVideo", s.HandlePlaceStreamMediaGetVideo) e.GET("/xrpc/place.stream.media.getVideoList", s.HandlePlaceStreamMediaGetVideoList) @@ -745,6 +746,24 @@ func (s *Server) HandlePlaceStreamMediaCreateUpload(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamMediaFinalizeLivestream(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamMediaFinalizeLivestream") + defer span.End() + + var body placestream.MediaFinalizeLivestream_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestream.MediaFinalizeLivestream_Output + var handleErr error + // func (s *Server) handlePlaceStreamMediaFinalizeLivestream(ctx context.Context,body *placestream.MediaFinalizeLivestream_Input) (*placestream.MediaFinalizeLivestream_Output, error) + out, handleErr = s.handlePlaceStreamMediaFinalizeLivestream(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandlePlaceStreamMediaGetUploadStatus(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamMediaGetUploadStatus") defer span.End() diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index 29a4b40f..a685beae 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -24,6 +24,7 @@ import ( var TaskNotification = "notification" var TaskChat = "chat" var TaskFinalizeLivestream = "finalize_livestream" +var TaskFinalizeLivestreamVOD = "finalize_livestream_vod" var TaskVODProcess = "vod_process" var TaskViewCountAggregate = "view_count_aggregate" @@ -69,6 +70,16 @@ type VODProcessTask struct { Location string `json:"location"` } +// FinalizeLivestreamVODTask is enqueued by the place.stream.media.finalizeLivestream +// procedure to turn a finished livestream's recorded MUXL objects into a VOD. +// UploadID is the synthetic Upload row the client polls (getUploadStatus) and +// then publishes (publishVideo), exactly as for a resumable upload. +type FinalizeLivestreamVODTask struct { + UploadID string `json:"uploadId"` + RepoDID string `json:"repoDID"` + LivestreamURI string `json:"livestreamURI"` +} + // ViewCountAggregateTask is the payload for one aggregation window. // Enqueued by every streamplace node at the configured interval; the // unique task key (built from WindowStart/End) ensures only one node's @@ -97,11 +108,13 @@ func (state *StatefulDB) ProcessQueue(ctx context.Context, vodConcurrency int) e return state.runQueueWorker(ctx, "queue_processor", nonVODTaskTypes) }) - // Dedicated VOD pool. + // Dedicated VOD pool. Live-to-VOD finalize also runs here: it's heavy + // I/O (a full read of the recorded stream to hash + index it) that would + // otherwise hog the single general worker and stall light tasks. for i := 0; i < vodConcurrency; i++ { workerID := fmt.Sprintf("vod_worker_%d", i) group.Go(func() error { - return state.runQueueWorker(ctx, workerID, []string{TaskVODProcess}) + return state.runQueueWorker(ctx, workerID, []string{TaskVODProcess, TaskFinalizeLivestreamVOD}) }) } @@ -144,6 +157,8 @@ func (state *StatefulDB) processTask(ctx context.Context, task *AppTask) error { return state.processFinalizeLivestreamTask(ctx, task) case TaskVODProcess: return state.processVODProcessTask(ctx, task) + case TaskFinalizeLivestreamVOD: + return state.processFinalizeLivestreamVODTask(ctx, task) case TaskViewCountAggregate: return state.processViewCountAggregateTask(ctx, task) default: @@ -192,6 +207,45 @@ func (state *StatefulDB) processVODProcessTask(ctx context.Context, task *AppTas return state.CompleteTask(ctx, task.ID) } +// LivestreamVODFinalizer concatenates a finished livestream's recorded MUXL +// objects into a content-addressed VOD blob, derives its sidecars, and +// publishes the origin + track records. Returns the resulting BDASL CID. Same +// function-pointer indirection as VODProcessor so pkg/statedb doesn't import +// the blob.Store/muxl-heavy pkg/vod. +type LivestreamVODFinalizer func(ctx context.Context, t FinalizeLivestreamVODTask) (cid string, err error) + +func (state *StatefulDB) SetLivestreamVODFinalizer(f LivestreamVODFinalizer) { + state.livestreamVODFinalizer = f +} + +func (state *StatefulDB) processFinalizeLivestreamVODTask(ctx context.Context, task *AppTask) error { + ctx = log.WithLogValues(ctx, "func", "processFinalizeLivestreamVODTask") + var t FinalizeLivestreamVODTask + if err := json.Unmarshal(task.Payload, &t); err != nil { + return err + } + if state.livestreamVODFinalizer == nil { + log.Warn(ctx, "no livestream VOD finalizer configured; dropping task", + "uploadId", t.UploadID, "did", t.RepoDID) + return state.CompleteTask(ctx, task.ID) + } + if err := state.SetUploadProcessing(ctx, t.UploadID); err != nil { + log.Warn(ctx, "failed to mark upload as processing", "uploadId", t.UploadID, "error", err) + } + cid, err := state.livestreamVODFinalizer(ctx, t) + if err != nil { + if ferr := state.SetUploadFailed(ctx, t.UploadID, err.Error()); ferr != nil { + log.Warn(ctx, "failed to mark upload as failed", "uploadId", t.UploadID, "error", ferr) + } + // Complete so it doesn't retry: most finalize failures are permanent + // (missing objects, unreadable bytes, no OAuth session). + _ = state.CompleteTask(ctx, task.ID) + return fmt.Errorf("finalize livestream VOD upload %s: %w", t.UploadID, err) + } + log.Log(ctx, "livestream VOD finalized", "uploadId", t.UploadID, "cid", cid) + return state.CompleteTask(ctx, task.ID) +} + // ViewCountAggregator runs the view-log → place.stream.media.viewCount // aggregation for one window. Same function-pointer indirection trick // as VODProcessor: pkg/statedb stays ignorant of pkg/viewlog (which diff --git a/pkg/statedb/s3_segment.go b/pkg/statedb/s3_segment.go index 8b18a166..3a11e6b0 100644 --- a/pkg/statedb/s3_segment.go +++ b/pkg/statedb/s3_segment.go @@ -9,17 +9,21 @@ import ( ) type S3Segment struct { - ID string `gorm:"column:id;primarykey"` - RepoDID string `gorm:"column:user_did;index;not null"` - Bucket string `gorm:"column:bucket;not null"` - Key string `gorm:"column:key;not null"` - URL string `gorm:"column:url"` - StartedAt time.Time `gorm:"column:started_at"` - CompletedAt *time.Time `gorm:"column:completed_at"` - Size int64 `gorm:"column:size"` - PartCount int32 `gorm:"column:part_count"` - CreatedAt time.Time `gorm:"column:created_at"` - UpdatedAt time.Time `gorm:"column:updated_at"` + ID string `gorm:"column:id;primarykey"` + RepoDID string `gorm:"column:user_did;index;not null"` + // LivestreamURI ties this object to the place.stream.livestream it was + // recorded for, so live-to-VOD finalize can enumerate exactly the objects + // of one stream. May be empty for objects started before the URI was known. + LivestreamURI string `gorm:"column:livestream_uri;index"` + Bucket string `gorm:"column:bucket;not null"` + Key string `gorm:"column:key;not null"` + URL string `gorm:"column:url"` + StartedAt time.Time `gorm:"column:started_at"` + CompletedAt *time.Time `gorm:"column:completed_at"` + Size int64 `gorm:"column:size"` + PartCount int32 `gorm:"column:part_count"` + CreatedAt time.Time `gorm:"column:created_at"` + UpdatedAt time.Time `gorm:"column:updated_at"` } func (s *S3Segment) TableName() string { @@ -28,17 +32,18 @@ func (s *S3Segment) TableName() string { // RecordStart inserts a new S3Segment row at the start of a multipart upload // and returns its ID. Implements s3.Recorder. -func (state *StatefulDB) RecordStart(ctx context.Context, repoDID, bucket, key string, started time.Time) (string, error) { +func (state *StatefulDB) RecordStart(ctx context.Context, repoDID, bucket, key, livestreamURI string, started time.Time) (string, error) { uu, err := uuid.NewV7() if err != nil { return "", err } seg := &S3Segment{ - ID: uu.String(), - RepoDID: repoDID, - Bucket: bucket, - Key: key, - StartedAt: started, + ID: uu.String(), + RepoDID: repoDID, + LivestreamURI: livestreamURI, + Bucket: bucket, + Key: key, + StartedAt: started, } if err := state.DB.WithContext(ctx).Create(seg).Error; err != nil { return "", err @@ -46,6 +51,23 @@ func (state *StatefulDB) RecordStart(ctx context.Context, repoDID, bucket, key s return seg.ID, nil } +// ListS3SegmentsForLivestream returns the completed S3 objects recorded for one +// livestream, ordered by StartedAt — i.e. the chronological byte order the +// objects must be concatenated in to reconstruct the stream. In-progress +// objects (CompletedAt == nil) are excluded; finalize runs after teardown, by +// which point every object of a finished stream has completed. +func (state *StatefulDB) ListS3SegmentsForLivestream(ctx context.Context, livestreamURI string) ([]S3Segment, error) { + var segs []S3Segment + err := state.DB.WithContext(ctx). + Where("livestream_uri = ? AND completed_at IS NOT NULL", livestreamURI). + Order("started_at ASC"). + Find(&segs).Error + if err != nil { + return nil, err + } + return segs, nil +} + // RecordComplete marks an S3 multipart upload as completed and records the // final part count and size. Implements s3.Recorder. func (state *StatefulDB) RecordComplete(ctx context.Context, id string, parts int32, size int64) error { diff --git a/pkg/statedb/statedb.go b/pkg/statedb/statedb.go index 32d7d18b..dc59bf1c 100644 --- a/pkg/statedb/statedb.go +++ b/pkg/statedb/statedb.go @@ -47,6 +47,10 @@ type StatefulDB struct { // SetViewCountAggregator at bootstrap so pkg/statedb doesn't have // to depend on the blob.Store-heavy pkg/viewlog. viewCountAggregator ViewCountAggregator + // livestreamVODFinalizer concatenates a finished livestream's recorded + // MUXL objects into a VOD. Installed via SetLivestreamVODFinalizer at + // bootstrap, same indirection as vodProcessor. + livestreamVODFinalizer LivestreamVODFinalizer } // list tables here so we can migrate them diff --git a/pkg/streamplace/mediafinalizeLivestream.go b/pkg/streamplace/mediafinalizeLivestream.go new file mode 100644 index 00000000..bf451c9c --- /dev/null +++ b/pkg/streamplace/mediafinalizeLivestream.go @@ -0,0 +1,33 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.media.finalizeLivestream + +package streamplace + +import ( + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// MediaFinalizeLivestream_Input is the input argument to a place.stream.media.finalizeLivestream call. +type MediaFinalizeLivestream_Input struct { + // livestream: AT-URI of the place.stream.livestream record to finalize into a VOD. Must belong to the authenticated user. + Livestream string `json:"livestream" cborgen:"livestream"` +} + +// MediaFinalizeLivestream_Output is the output of a place.stream.media.finalizeLivestream call. +type MediaFinalizeLivestream_Output struct { + // uploadId: Identifier for the finalize job. Poll place.stream.media.getUploadStatus with it; once status is 'done', create the video with place.stream.media.publishVideo. + UploadId string `json:"uploadId" cborgen:"uploadId"` +} + +// MediaFinalizeLivestream calls the XRPC method "place.stream.media.finalizeLivestream". +func MediaFinalizeLivestream(ctx context.Context, c lexutil.LexClient, input *MediaFinalizeLivestream_Input) (*MediaFinalizeLivestream_Output, error) { + var out MediaFinalizeLivestream_Output + if err := c.LexDo(ctx, lexutil.Procedure, "application/json", "place.stream.media.finalizeLivestream", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/vod/finalize_livestream.go b/pkg/vod/finalize_livestream.go new file mode 100644 index 00000000..467c58bc --- /dev/null +++ b/pkg/vod/finalize_livestream.go @@ -0,0 +1,335 @@ +package vod + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "strings" + + "go.opentelemetry.io/otel/attribute" + "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" + "stream.place/streamplace/pkg/media" + "stream.place/streamplace/pkg/muxl" + s3pkg "stream.place/streamplace/pkg/s3" + "stream.place/streamplace/pkg/statedb" +) + +// FinalizeInput drives FinalizeLivestreamVOD. It's resolved by the bootstrap +// from the FinalizeLivestreamVODTask payload (plus a model lookup for the +// signing key), so this package needn't import pkg/model. +type FinalizeInput struct { + // UploadID is the synthetic Upload row created by the finalizeLivestream + // procedure; the getUploadStatus/publishVideo client contract keys off it. + UploadID string + // RepoDID is the streamer whose repo the media.track records land in. + RepoDID string + // LivestreamURI selects which S3Segment objects belong to this stream. + LivestreamURI string + // SigningKey is the did:key whose private half C2PA-signed the live + // segments (the streamer's place.stream.key). Recorded on each + // place.stream.media.track so playback can verify signatures. + SigningKey string +} + +// FinalizeLivestreamVOD turns a finished livestream's recorded MUXL objects +// into a VOD, reusing the upload pipeline's tail half. The live segments are +// already C2PA-signed canonical MUXL, so — unlike ProcessVOD — there is no +// gstreamer remux and no re-signing: we concatenate the objects under one +// synthesized init header into a content-addressed blob, derive the playback +// sidecars from it with `muxl unwrap` (exactly as VOD transfer does), and +// publish the same origin + track records. +// +// The content blob is shaped [init][segments…] — byte-structurally identical +// to an uploaded VOD — so the metafile offset math is the proven path and a +// live-derived VOD is indistinguishable from an uploaded one downstream +// (playback, transfer, view-count). The bulk of the bytes are assembled +// server-side via S3 UploadPartCopy; only the ~5 MB header part is uploaded. +func FinalizeLivestreamVOD(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, store blob.Store, in FinalizeInput) (string, error) { + ctx = log.WithLogValues(ctx, "func", "FinalizeLivestreamVOD", "uploadId", in.UploadID, "did", in.RepoDID, "livestream", in.LivestreamURI) + ctx, span := vodTracer.Start(ctx, "vod.FinalizeLivestreamVOD", trace.WithAttributes( + attribute.String("upload_id", in.UploadID), + attribute.String("did", in.RepoDID), + attribute.String("livestream", in.LivestreamURI), + )) + defer span.End() + + segs, err := state.ListS3SegmentsForLivestream(ctx, in.LivestreamURI) + if err != nil { + recordErr(span, "list_segments", err) + return "", fmt.Errorf("list s3 segments: %w", err) + } + if len(segs) == 0 { + err := fmt.Errorf("no recorded S3 segments for livestream %s", in.LivestreamURI) + recordErr(span, "list_segments", err) + return "", err + } + keys := make([]string, len(segs)) + for i, s := range segs { + keys[i] = s.Key + } + 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]) + if err != nil { + recordErr(span, "capture_init", err) + return "", fmt.Errorf("capture init header: %w", err) + } + span.SetAttributes(attribute.Int("header_bytes", len(header))) + + // 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) + if err != nil { + recordErr(span, "build_metafile", err) + return "", err + } + contentKey := BlobsPrefix + cid + ".mp4" + span.SetAttributes( + attribute.String("cid", cid), + attribute.Int64("size_bytes", size), + attribute.String("content_key", contentKey), + ) + + // 3. Assemble the content blob, mostly server-side. + if err := runVODStage(ctx, "concat_assemble", func(ctx context.Context) error { + return assembleContentBlob(ctx, store, header, 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) + }); 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. + 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, + mimeType: "video/mp4", + probe: probe, + signingKey: in.SigningKey, + }) + }); err != nil { + recordErr(span, stagePublish, err) + return "", fmt.Errorf("publish records: %w", err) + } + + span.SetStatus(codes.Ok, "") + log.Log(ctx, "livestream VOD finalized", "cid", cid, "size", size, "objects", len(keys), "duration_ms", probe.DurationMS) + return cid, nil +} + +// captureInitHeader runs `muxl unwrap` over just the first object and returns +// the bytes of the single init event it emits, cancelling as soon as it has +// them so we don't read the whole object. The returned header is a valid +// ftyp+moov for the stream's catalog. +func captureInitHeader(ctx context.Context, store blob.Store, key string) ([]byte, error) { + r, err := store.Open(ctx, key) + if err != nil { + return nil, fmt.Errorf("open first object %s: %w", key, err) + } + defer r.Close() + + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + eventCh := make(chan *muxl.MuxlEvent, 16) + producerErr := make(chan error, 1) + go func() { + producerErr <- muxl.RunMuxlUnwrapEvents(ctx, io.NewSectionReader(r, 0, r.Size()), eventCh) + close(eventCh) + }() + + var header []byte + for ev := range eventCh { + if header == nil && ev.Type == "init" { + header = append([]byte(nil), ev.Data...) + cancel() // we have the header; let the producer unwind, keep draining + } + } + // RunMuxlUnwrapEvents returns nil on cancellation, so a non-nil error here + // is a real parse failure, not our early stop. + if perr := <-producerErr; perr != nil && header == nil { + return nil, fmt.Errorf("muxl unwrap (init): %w", perr) + } + if header == nil { + return nil, errors.New("no init event emitted by first object") + } + return header, 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 +// full rewrite, as does any non-S3 store. +func assembleContentBlob(ctx context.Context, store blob.Store, header []byte, keys []string, contentKey string) error { + if s3store, ok := store.(*blob.S3Store); ok { + err := s3pkg.ConcatWithHeader(ctx, s3store.Client(), s3store.Bucket(), header, keys, contentKey, "video/mp4") + if err == nil { + return nil + } + if !errors.Is(err, s3pkg.ErrConcatPartTooSmall) { + return err + } + log.Warn(ctx, "s3 concat fell back to rewrite (object below 5MB part floor)", "key", contentKey) + } + return assembleViaWriter(ctx, store, header, keys, contentKey) +} + +// assembleViaWriter streams header ++ objects through a store writer. The +// universal fallback: correct for any blob.Store, but re-uploads every byte. +func assembleViaWriter(ctx context.Context, store blob.Store, header []byte, keys []string, contentKey string) error { + w, err := store.NewWriter(ctx, contentKey, "video/mp4") + if err != nil { + return fmt.Errorf("open content writer: %w", err) + } + defer w.Close() + if _, err := w.Write(header); err != nil { + return fmt.Errorf("write header: %w", err) + } + for _, key := range keys { + r, err := store.Open(ctx, key) + if err != nil { + return fmt.Errorf("open object %s: %w", key, err) + } + _, err = io.Copy(w, io.NewSectionReader(r, 0, r.Size())) + _ = r.Close() + if err != nil { + return fmt.Errorf("copy object %s: %w", key, err) + } + } + return w.Complete() +} + +// metafileToVODResult derives the probe metadata publishTrack needs from the +// regenerated metafile, standing in for the gstreamer probe that an upload +// would have produced. publishTrack hardcodes the video codec to h264 and only +// uses width/height (+ optional framerate), so we leave framerate unset. +func metafileToVODResult(meta *Metafile) media.VODResult { + var res media.VODResult + for _, t := range meta.Tracks { + switch t.Type { + case "video": + res.Video = &media.VODVideoTrack{ + Codec: t.Codec, + Width: int(t.Width), + Height: int(t.Height), + } + case "audio": + res.Audio = &media.VODAudioTrack{ + Codec: audioProbeCodec(t.Codec), + Rate: int(t.SampleRate), + Channels: int(t.Channels), + } + } + if d := metafileTrackDurationMS(t); d > res.DurationMS { + res.DurationMS = d + } + } + return res +} + +// metafileTrackDurationMS sums a track's segment durations (in its timescale) +// and converts to milliseconds. +func metafileTrackDurationMS(t MetafileTrack) int64 { + if t.Timescale == 0 { + return 0 + } + var ticks uint64 + for _, s := range t.Segments { + ticks += s.DurationTicks + } + return int64(ticks * 1000 / uint64(t.Timescale)) +} + +// audioProbeCodec maps a metafile (CMAF catalog) audio codec name to the +// gstreamer-caps-style value audioCodecForLexicon (publish.go) expects, so the +// lexicon enum comes out right ("opus" vs "aac"). +func audioProbeCodec(metafileCodec string) string { + if strings.Contains(strings.ToLower(metafileCodec), "opus") { + return "x-opus" + } + return "mpeg" // audioCodecForLexicon maps "mpeg" -> "aac" +} diff --git a/pkg/vod/finalize_livestream_test.go b/pkg/vod/finalize_livestream_test.go new file mode 100644 index 00000000..f1b8a5bc --- /dev/null +++ b/pkg/vod/finalize_livestream_test.go @@ -0,0 +1,158 @@ +package vod + +import ( + "bytes" + "context" + "os" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/bdasl" + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/log" +) + +// TestFinalizeHeaderConcat is the gating test for live-to-VOD finalize: it +// proves that prepending the init `muxl unwrap` synthesizes for a stream's bare +// signed segments produces a content blob whose regenerated metafile is +// playback-correct — same per-track init CIDs and same per-segment sizes / +// durations as the processing pipeline's, with self-consistent (contiguous, +// header-offset) byte ranges. +// +// This is the assumption FinalizeLivestreamVOD rests on: the live S3 objects +// are bare segments, and finalize prepends one synthesized header so the blob +// is shaped [init][segments…] like an uploaded VOD. If muxl's bare-input init +// synthesis ever diverged from what the metafile builder assumes for a leading +// init, the byte ranges would be wrong and this fails. +// +// Like the other pkg/vod tests it needs gstreamer (warmGST/streamThroughMuxl), +// so it runs in the cgo test containers, not a bare checkout. +func TestFinalizeHeaderConcat(t *testing.T) { + warmGST() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + ctx = log.WithLogValues(ctx, "test", "TestFinalizeHeaderConcat") + + fixture, err := os.ReadFile(getFixture("5sec.mp4")) + require.NoError(t, err) + + signer, err := newUploadSigner(time.Now()) + require.NoError(t, err) + + // --- Ground truth: the processing pipeline's blob + metafile. ---------- + storeA, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + + out := &bytes.Buffer{} + hasher := bdasl.NewWriter() + dst := teeWriter{hasher, out} + mbA := newMetafileBuilder(ctx, storeA) + _, err = streamThroughMuxl(ctx, bytes.NewReader(fixture), int64(len(fixture)), dst, mbA, signer.SignerInput) + require.NoError(t, err) + groundBlob := out.Bytes() + metaA := mbA.Finalize(hasher.CID(), int64(len(groundBlob))) + require.NotEmpty(t, metaA.Tracks) + + // The processing blob is [init][segments…]; the live S3 objects would be + // the bare segments (no leading init). Strip the init to simulate them: it + // occupies [0, firstSegmentOffset). + initLen := minFirstOffset(t, metaA) + require.Greater(t, initLen, int64(0)) + bareSegs := groundBlob[initLen:] + + // --- Finalize path: capture header, prepend, regenerate. --------------- + storeB, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + + // Write the bare segments as the single recorded "object". + objKey := "live/obj0.m4s" + wObj, err := storeB.NewWriter(ctx, objKey, "video/iso.segment") + require.NoError(t, err) + _, err = wObj.Write(bareSegs) + require.NoError(t, err) + require.NoError(t, wObj.Complete()) + + // captureInitHeader is the real finalize helper: unwrap the bare object and + // take the init event it emits. + header, err := captureInitHeader(ctx, storeB, objKey) + require.NoError(t, err) + require.NotEmpty(t, header) + + newBlob := append(append([]byte(nil), header...), bareSegs...) + cidB := func() string { + h := bdasl.NewWriter() + _, _ = h.Write(newBlob) + return h.CID() + }() + contentKey := BlobsPrefix + cidB + ".mp4" + wBlob, err := storeB.NewWriter(ctx, contentKey, "video/mp4") + require.NoError(t, err) + _, err = wBlob.Write(newBlob) + require.NoError(t, err) + require.NoError(t, wBlob.Complete()) + + metaB, initCIDs, err := regenerateSidecars(ctx, storeB, cidB, int64(len(newBlob))) + require.NoError(t, err) + require.NotEmpty(t, initCIDs) + + // --- Correctness assertions. ------------------------------------------- + require.Equal(t, len(metaA.Tracks), len(metaB.Tracks), "track count") + for tid, ta := range metaA.Tracks { + tb, ok := metaB.Tracks[tid] + require.Truef(t, ok, "track %s missing from regenerated metafile", tid) + + // Per-track identity must match: same init blob, codec, geometry. + require.Equalf(t, ta.InitCID, tb.InitCID, "track %s initCID", tid) + require.Equalf(t, ta.Codec, tb.Codec, "track %s codec", tid) + require.Equalf(t, ta.Type, tb.Type, "track %s type", tid) + require.Equalf(t, ta.Timescale, tb.Timescale, "track %s timescale", tid) + require.Equalf(t, ta.Width, tb.Width, "track %s width", tid) + require.Equalf(t, ta.Height, tb.Height, "track %s height", tid) + require.Equalf(t, ta.Channels, tb.Channels, "track %s channels", tid) + require.Equalf(t, ta.SampleRate, tb.SampleRate, "track %s sampleRate", tid) + + // Per-segment sizes/durations must match exactly; offsets may differ + // only if the synthesized header length differs from the pipeline init. + require.Equalf(t, len(ta.Segments), len(tb.Segments), "track %s segment count", tid) + for i := range ta.Segments { + require.Equalf(t, ta.Segments[i].Size, tb.Segments[i].Size, "track %s seg %d size", tid, i) + require.Equalf(t, ta.Segments[i].DurationTicks, tb.Segments[i].DurationTicks, "track %s seg %d duration", tid, i) + require.Equalf(t, ta.Segments[i].SampleCount, tb.Segments[i].SampleCount, "track %s seg %d sampleCount", tid, i) + } + } + + // The regenerated metafile must be self-consistent with the actual blob: + // the earliest segment starts exactly after the header, and each track's + // segments tile contiguously. + require.Equal(t, int64(len(header)), minFirstOffset(t, metaB), "first segment must start right after the header") + for tid, tb := range metaB.Tracks { + for i := 1; i < len(tb.Segments); i++ { + prev := tb.Segments[i-1] + require.Equalf(t, prev.Offset+prev.Size, tb.Segments[i].Offset, + "track %s seg %d not contiguous", tid, i) + } + } + + // Sanity: the prepended-header blob round-trips to its own CID, and the + // header length matches the pipeline's init length (so a live VOD is + // byte-structurally identical to an uploaded one). + require.Equal(t, initLen, int64(len(header)), "synthesized header length should match the pipeline init length") +} + +// minFirstOffset returns the smallest first-segment offset across all tracks — +// i.e. where the first segment byte begins, which equals the leading init's +// length in a [init][segments…] blob. +func minFirstOffset(t *testing.T, m *Metafile) int64 { + t.Helper() + var min int64 = -1 + for _, tr := range m.Tracks { + require.NotEmpty(t, tr.Segments) + off := tr.Segments[0].Offset + if min < 0 || off < min { + min = off + } + } + return min +}