From e8a011cb82c4ab2faccfd9cae388cac1de8feb00 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 26 Jun 2026 18:21:44 -0700 Subject: [PATCH] draft-vods: defer media.track publication to publishDraft time place.stream.media.track records were published at processing time (in publishRecords), which leaked half-published content: the tracks went live immediately while the video record wasn't created until the user published the draft. With drafts, that gap can be long. Defer track creation to publishDraft. - publishRecords stops calling publishTrack. It still publishes place.stream.media.origin (server-repo blob availability attestation) and now stores the track-publishing inputs on the Upload row: signing_key, probe_json (video/audio codec/dims/fps/rate/channels), and blob_size. SetUploadProcessed signature updated accordingly. - markDraftReadyFromUpload: the draft reaches 'ready' with source=nil (no track refs exist yet) and durationMs set (known at processing time). - PublishDraft publishes the tracks at publish time via the new publishTracksFromUpload: reads the probe+signingKey+size+CID from the tied Upload row, calls publishTrack per A/V stream, and builds the video's source from the fresh strongRefs. Falls back to a carried-over source if the draft has no tied upload (legacy). - Upload row gains SigningKey, ProbeJSON, BlobSize columns. TrackURIs is vestigial (left for the legacy publishVideo path). Tests updated: the lifecycle test now asserts source is nil at ready time and the Upload row carries the deferred inputs. build + tests + lint clean. Co-authored-by: Claude Opus 4.8 --- pkg/statedb/draft_video.go | 15 ++- pkg/statedb/queue_processor_draft_test.go | 20 ++-- pkg/statedb/upload.go | 34 +++++- pkg/vod/publish.go | 127 +++++++++++++++------- pkg/vod/publish_draft.go | 84 ++++++++++++-- 5 files changed, 214 insertions(+), 66 deletions(-) diff --git a/pkg/statedb/draft_video.go b/pkg/statedb/draft_video.go index 0729c95c..861185fc 100644 --- a/pkg/statedb/draft_video.go +++ b/pkg/statedb/draft_video.go @@ -236,11 +236,13 @@ func (state *StatefulDB) SetDraftReady(ctx context.Context, originUploadID strin } // markDraftReadyFromUpload re-reads a (now-finished) Upload row and flips its -// tied draft to 'ready', filling source/durationMs from the row and content_cid -// for later thumbnail generation. Called by the queue processors after the VOD -// processor / livestream finalizer returns, since those call SetUploadProcessed -// internally and return only a cid — the processor's results land on the Upload -// row, which we re-read here. A missing draft is a no-op. +// tied draft to 'ready', filling durationMs + content_cid. The draft's source +// is left empty: track records are deferred to publishDraft time, so there +// are no track refs to populate it with at ready time. Called by the queue +// processors after the VOD processor / livestream finalizer returns, since +// those call SetUploadProcessed internally and return only a cid — the +// processor's results land on the Upload row, which we re-read here. A missing +// draft is a no-op. func (state *StatefulDB) markDraftReadyFromUpload(ctx context.Context, uploadID string) error { upload, err := state.GetUpload(ctx, uploadID) if err != nil { @@ -249,6 +251,9 @@ func (state *StatefulDB) markDraftReadyFromUpload(ctx context.Context, uploadID if upload == nil { return nil } + // source stays nil until PublishDraft publishes the tracks; durationMs is + // known at processing time. (draftSourceFromTrackURIs returns nil for the + // now-empty TrackURIs, so the legacy path is unchanged.) source, err := draftSourceFromTrackURIs(upload.TrackURIs) if err != nil { return err diff --git a/pkg/statedb/queue_processor_draft_test.go b/pkg/statedb/queue_processor_draft_test.go index 7df6f15d..4bc6f203 100644 --- a/pkg/statedb/queue_processor_draft_test.go +++ b/pkg/statedb/queue_processor_draft_test.go @@ -36,9 +36,11 @@ func TestDraftLifecycleThroughVODProcessor(t *testing.T) { // Fake processor: simulates vod.ProcessVOD's internal SetUploadProcessed // (writes the finished fields onto the Upload row) then returns a cid. - trackURIs := `[{"uri":"at://did:plc:trackhost/place.stream.media.track/t1","cid":"bafytrack1"}]` + // Tracks are deferred to publishDraft time now, so SetUploadProcessed + // stores the probe + signingKey (not track refs). + probeJSON := `{"durationMs":98765,"video":{"codec":"h264","width":1280,"height":720,"fpsNum":30,"fpsDen":1}}` state.SetVODProcessor(func(ctx context.Context, task VODProcessTask) (string, error) { - require.NoError(t, state.SetUploadProcessed(ctx, task.UploadID, trackURIs, 98765, "muxlcid-lifecycle")) + require.NoError(t, state.SetUploadProcessed(ctx, task.UploadID, 98765, "muxlcid-lifecycle", "did:key:signing", probeJSON, 123456)) return "muxlcid-lifecycle", nil }) @@ -50,7 +52,9 @@ func TestDraftLifecycleThroughVODProcessor(t *testing.T) { })} _ = state.processVODProcessTask(ctx, task) - // The draft must now be 'ready' and carry the upload's values. + // The draft must now be 'ready'. durationMs + content_cid are filled + // from the upload row; source is nil (tracks are published at + // publishDraft time, not at ready time). dv, err := state.GetDraftByUpload(ctx, "up-lifecycle") require.NoError(t, err) require.NotNil(t, dv) @@ -60,10 +64,12 @@ func TestDraftLifecycleThroughVODProcessor(t *testing.T) { require.Equal(t, "ready", rec.Status) require.NotNil(t, rec.DurationMs) require.Equal(t, int64(98765), *rec.DurationMs) - require.NotNil(t, rec.Source) - require.NotNil(t, rec.Source.MediaDefs_SourceTracks) - require.Len(t, rec.Source.MediaDefs_SourceTracks.Tracks, 1) - require.Equal(t, "at://did:plc:trackhost/place.stream.media.track/t1", rec.Source.MediaDefs_SourceTracks.Tracks[0].Uri) + require.Nil(t, rec.Source, "source must be nil at ready time (tracks deferred to publish)") + // The Upload row carries the deferred track-publishing inputs. + up, err := state.GetUpload(ctx, "up-lifecycle") + require.NoError(t, err) + require.Equal(t, "did:key:signing", up.SigningKey) + require.Contains(t, up.ProbeJSON, "h264") }) } diff --git a/pkg/statedb/upload.go b/pkg/statedb/upload.go index 0c45ac48..921aab9a 100644 --- a/pkg/statedb/upload.go +++ b/pkg/statedb/upload.go @@ -31,9 +31,11 @@ type Upload struct { ProcessingStatus string `gorm:"column:processing_status"` ProcessingError string `gorm:"column:processing_error"` ProcessingProgress int `gorm:"column:processing_progress;default:0"` - // TrackURIs is a JSON array of {"uri":"at://...","cid":"..."} objects - // populated once the track records are published and the video is ready - // for the client to create a place.stream.video record. + // TrackURIs is a JSON array of {"uri":"at://...","cid":"..."} objects. + // Vestigial since track publication was deferred to publishDraft time: + // new uploads leave this empty and PublishDraft publishes the tracks on + // demand. Retained so existing rows / the legacy publishVideo path still + // read it. TrackURIs string `gorm:"column:track_uris"` DurationMS int64 `gorm:"column:duration_ms"` // ContentCID is the BDASL CID of the processed fMP4 blob. Stored so the @@ -41,6 +43,20 @@ type Upload struct { // publishVideo's thumbnail generation) without re-deriving it from the // published track records. ContentCID string `gorm:"column:content_cid"` + // SigningKey is the did:key whose ephemeral private half C2PA-signed the + // segments. Stored at processing time so PublishDraft can publish the + // place.stream.media.track records (which carry it) at publish time. + SigningKey string `gorm:"column:signing_key"` + // ProbeJSON is the gstreamer probe metadata (video/audio codec, dims, + // fps, rate, channels) serialized as JSON. Stored at processing time so + // PublishDraft can publish the track records (deferred from processing) + // without re-probing the blob. + ProbeJSON string `gorm:"column:probe_json"` + // BlobSize is the byte size of the processed MUXL content blob (distinct + // from Size, the raw upload size). Stored at processing time so + // PublishDraft can populate the track records' size field without + // re-statting the blob. + BlobSize int64 `gorm:"column:blob_size"` } func (Upload) TableName() string { @@ -88,15 +104,23 @@ func (state *StatefulDB) SetUploadProgress(ctx context.Context, id string, progr Update("processing_progress", progress).Error } -func (state *StatefulDB) SetUploadProcessed(ctx context.Context, id string, trackURIsJSON string, durationMS int64, contentCID string) error { +// SetUploadProcessed marks an upload done and stores the processing results +// the publishDraft path needs: duration, the content blob's CID, the C2PA +// signing key, and the gstreamer probe metadata (JSON). Track records are +// NOT published here — they're deferred to publishDraft time so half- +// published tracks don't go live before the video record. trackURIs is left +// empty (vestigial; the legacy publishVideo path may still set it). +func (state *StatefulDB) SetUploadProcessed(ctx context.Context, id string, durationMS int64, contentCID, signingKey, probeJSON string, blobSize int64) error { return state.DB.WithContext(ctx).Model(&Upload{}). Where("id = ?", id). Updates(map[string]any{ "processing_status": "done", "processing_progress": 100, - "track_uris": trackURIsJSON, "duration_ms": durationMS, "content_cid": contentCID, + "signing_key": signingKey, + "probe_json": probeJSON, + "blob_size": blobSize, }).Error } diff --git a/pkg/vod/publish.go b/pkg/vod/publish.go index 57cc7090..b24d8663 100644 --- a/pkg/vod/publish.go +++ b/pkg/vod/publish.go @@ -48,18 +48,19 @@ type publishParams struct { signingKey string } -// publishRecords does the post-processing record publish: +// publishRecords does the post-processing record publish. With tracks +// deferred to publishDraft time, this now: // -// 1. place.stream.media.origin in the SERVER's repo (we attest that -// this blob is fetchable from us). Idempotent: rkey is the CID. -// 2. place.stream.media.track in the USER's repo, one per A/V track, -// via the user's stored OAuth session. -// 3. Stores the resulting track URIs + duration on the Upload row so -// the client can poll getUploadStatus and create the -// place.stream.video record itself (with full metadata) via Publish. +// 1. place.stream.media.origin in the SERVER's repo (we attest that this +// blob is fetchable from us). Idempotent: rkey is the CID. +// 2. Stores the probe metadata + signing key + duration + CID on the Upload +// row so publishDraft can publish the place.stream.media.track records +// (and build the video's source) at publish time. // -// The video record is intentionally NOT created here — the client -// controls when it becomes visible and supplies the metadata. +// Track records are intentionally NOT published here: with drafts, the video +// record isn't visible until the user publishes, so publishing tracks at +// processing time would leak half-published content. publishDraft publishes +// them when the video record is created. func publishRecords(ctx context.Context, p publishParams) error { ctx, span := vodTracer.Start(ctx, "vod.publishRecords", trace.WithAttributes( attribute.String("cid", p.cid), @@ -73,48 +74,92 @@ func publishRecords(ctx context.Context, p publishParams) error { return fmt.Errorf("publish origin: %w", err) } - client, err := getUserXRPCClient(ctx, p.state, p.in.RepoDID) + // Serialize the probe so publishDraft can publish the track records later + // without re-probing the blob. + probeJSON, err := marshalProbe(p.probe) if err != nil { span.RecordError(err) - span.SetStatus(codes.Error, "get_client") - return fmt.Errorf("get user xrpc client: %w", err) + span.SetStatus(codes.Error, "marshal_probe") + return fmt.Errorf("marshal probe: %w", err) } - var trackRefs []trackRefJSON - if p.probe.Video != nil { - ref, err := publishTrack(ctx, client, p.in.RepoDID, p.cid, p.size, p.probe.DurationMS, "1", "video", p.signingKey, p.probe.Video, nil) - if err != nil { - span.RecordError(err) - span.SetStatus(codes.Error, "video_track") - return fmt.Errorf("publish video track: %w", err) + if err := p.state.SetUploadProcessed(ctx, p.in.UploadID, p.probe.DurationMS, p.cid, p.signingKey, probeJSON, p.size); err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, "store_results") + return fmt.Errorf("store processing results: %w", err) + } + + span.SetStatus(codes.Ok, "") + log.Log(ctx, "stored processing results (tracks deferred to publish)", + "uploadId", p.in.UploadID, "cid", p.cid, "duration_ms", p.probe.DurationMS) + return nil +} + +// probeJSONShape mirrors media.VODResult's track fields, for serialization to +// the Upload row's probe_json column. Only the fields publishTrack consumes. +type probeJSONShape struct { + DurationMS int64 `json:"durationMs"` + Video *videoProbeJSON `json:"video,omitempty"` + Audio *audioProbeJSON `json:"audio,omitempty"` +} +type videoProbeJSON struct { + Codec string `json:"codec"` + Width int `json:"width"` + Height int `json:"height"` + FPSNum int `json:"fpsNum"` + FPSDen int `json:"fpsDen"` +} +type audioProbeJSON struct { + Codec string `json:"codec"` + Rate int `json:"rate"` + Channels int `json:"channels"` + MPEGVersion int `json:"mpegVersion"` +} + +func marshalProbe(p media.VODResult) (string, error) { + out := probeJSONShape{DurationMS: p.DurationMS} + if p.Video != nil { + out.Video = &videoProbeJSON{ + Codec: p.Video.Codec, Width: p.Video.Width, Height: p.Video.Height, + FPSNum: p.Video.FPSNum, FPSDen: p.Video.FPSDen, } - trackRefs = append(trackRefs, trackRefJSON{URI: ref.Uri, CID: ref.Cid}) } - if p.probe.Audio != nil { - ref, err := publishTrack(ctx, client, p.in.RepoDID, p.cid, p.size, p.probe.DurationMS, "2", "audio", p.signingKey, nil, p.probe.Audio) - if err != nil { - span.RecordError(err) - span.SetStatus(codes.Error, "audio_track") - return fmt.Errorf("publish audio track: %w", err) + if p.Audio != nil { + out.Audio = &audioProbeJSON{ + Codec: p.Audio.Codec, Rate: p.Audio.Rate, Channels: p.Audio.Channels, + MPEGVersion: p.Audio.MPEGVersion, } - trackRefs = append(trackRefs, trackRefJSON{URI: ref.Uri, CID: ref.Cid}) } - - trackURIsJSON, err := json.Marshal(trackRefs) + b, err := json.Marshal(out) if err != nil { - span.RecordError(err) - span.SetStatus(codes.Error, "marshal_tracks") - return fmt.Errorf("marshal track refs: %w", err) - } - if err := p.state.SetUploadProcessed(ctx, p.in.UploadID, string(trackURIsJSON), p.probe.DurationMS, p.cid); err != nil { - span.RecordError(err) - span.SetStatus(codes.Error, "store_tracks") - return fmt.Errorf("store track refs: %w", err) + return "", err } + return string(b), nil +} - span.SetAttributes(attribute.Int("track_count", len(trackRefs))) - span.SetStatus(codes.Ok, "") - return nil +// unmarshalProbe reverses marshalProbe. +func unmarshalProbe(s string) (media.VODResult, error) { + if s == "" { + return media.VODResult{}, nil + } + var pjs probeJSONShape + if err := json.Unmarshal([]byte(s), &pjs); err != nil { + return media.VODResult{}, fmt.Errorf("unmarshal probe: %w", err) + } + res := media.VODResult{DurationMS: pjs.DurationMS} + if pjs.Video != nil { + res.Video = &media.VODVideoTrack{ + Codec: pjs.Video.Codec, Width: pjs.Video.Width, Height: pjs.Video.Height, + FPSNum: pjs.Video.FPSNum, FPSDen: pjs.Video.FPSDen, + } + } + if pjs.Audio != nil { + res.Audio = &media.VODAudioTrack{ + Codec: pjs.Audio.Codec, Rate: pjs.Audio.Rate, Channels: pjs.Audio.Channels, + MPEGVersion: pjs.Audio.MPEGVersion, + } + } + return res, nil } // publishOrigin attests that this server has the blob with the given diff --git a/pkg/vod/publish_draft.go b/pkg/vod/publish_draft.go index 722422d1..04b8ea81 100644 --- a/pkg/vod/publish_draft.go +++ b/pkg/vod/publish_draft.go @@ -84,8 +84,32 @@ func PublishDraft(ctx context.Context, state *statedb.StatefulDB, store blob.Sto video.ContentRights = rec.ContentRights video.Thumb = rec.Thumb - // Carry over the source union (the published track refs). - if rec.Source != nil && rec.Source.MediaDefs_SourceTracks != nil { + // Publish the place.stream.media.track records now (deferred from + // processing time so they don't leak before the video record). Read the + // probe + signing key from the tied Upload row, publish one track per + // A/V stream, and build the video's source from the fresh strongRefs. + // If there's no Upload row (e.g. a draft whose upload predates this + // change, or a draft created outside the upload flow), fall back to + // carrying over any source the draft already carries. + client, err := getUserXRPCClient(ctx, state, did) + if err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, "get_client") + return "", "", fmt.Errorf("get user xrpc client: %w", err) + } + sourceTracks, terr := publishTracksFromUpload(ctx, state, client, did, dv) + if terr != nil { + span.RecordError(terr) + span.SetStatus(codes.Error, "publish_tracks") + return "", "", fmt.Errorf("publish tracks: %w", terr) + } + if sourceTracks != nil { + video.Source = &streamplace.Video_Source{ + MediaDefs_SourceTracks: sourceTracks, + } + } else if rec.Source != nil && rec.Source.MediaDefs_SourceTracks != nil { + // Legacy fallback: a draft that already carries source (e.g. published + // via the old at-ready-time path). Carry it over. video.Source = &streamplace.Video_Source{ MediaDefs_SourceTracks: rec.Source.MediaDefs_SourceTracks, } @@ -95,12 +119,6 @@ func PublishDraft(ctx context.Context, state *statedb.StatefulDB, store blob.Sto } } - client, err := getUserXRPCClient(ctx, state, did) - if err != nil { - span.RecordError(err) - return "", "", fmt.Errorf("get user xrpc client: %w", err) - } - // Reuse the draft's tid as the published video's rkey. The draft is // deleted after a successful publish, so there's no collision, and sharing // the tid lets a single /upload/video/ route render either a draft @@ -172,6 +190,56 @@ func draftTID(atsURI string) (string, error) { return atsURI[idx+1:], nil } +// publishTracksFromUpload publishes the place.stream.media.track records for a +// draft's tied upload, deferred from processing time to publish time. It reads +// the probe metadata + signing key + blob size + CID from the Upload row and +// calls publishTrack for each A/V stream, returning the fresh strongRefs to use +// as the video record's source. Returns (nil, nil) if the draft has no tied +// upload (so the caller can fall back to a carried-over source). +func publishTracksFromUpload(ctx context.Context, state *statedb.StatefulDB, client XRPCClient, did string, dv *statedb.DraftVideo) (*streamplace.MediaDefs_SourceTracks, error) { + if dv.OriginUploadID == "" { + return nil, nil + } + upload, err := state.GetUpload(ctx, dv.OriginUploadID) + if err != nil { + return nil, fmt.Errorf("get upload: %w", err) + } + if upload == nil { + return nil, nil + } + probe, err := unmarshalProbe(upload.ProbeJSON) + if err != nil { + return nil, err + } + if upload.ContentCID == "" { + return nil, fmt.Errorf("upload %s has no content_cid", upload.ID) + } + var tracks []*comatproto.RepoStrongRef + if probe.Video != nil { + ref, err := publishTrack(ctx, client, did, upload.ContentCID, upload.BlobSize, probe.DurationMS, "1", "video", upload.SigningKey, probe.Video, nil) + if err != nil { + return nil, fmt.Errorf("publish video track: %w", err) + } + tracks = append(tracks, ref) + } + if probe.Audio != nil { + ref, err := publishTrack(ctx, client, did, upload.ContentCID, upload.BlobSize, probe.DurationMS, "2", "audio", upload.SigningKey, nil, probe.Audio) + if err != nil { + return nil, fmt.Errorf("publish audio track: %w", err) + } + tracks = append(tracks, ref) + } + if len(tracks) == 0 { + return nil, nil + } + log.Log(ctx, "published media.track records (deferred to publish)", + "uploadId", upload.ID, "cid", upload.ContentCID, "tracks", len(tracks)) + return &streamplace.MediaDefs_SourceTracks{ + LexiconTypeID: "place.stream.media.defs#sourceTracks", + Tracks: tracks, + }, nil +} + func derefInt64(p *int64) int64 { if p == nil { return 0 -- 2.51.2