diff --git a/pkg/spxrpc/place_stream_media_finalizelivestream.go b/pkg/spxrpc/place_stream_media_finalizelivestream.go index 1e2c9feb..7bebfb35 100644 --- a/pkg/spxrpc/place_stream_media_finalizelivestream.go +++ b/pkg/spxrpc/place_stream_media_finalizelivestream.go @@ -67,5 +67,14 @@ func (s *Server) handlePlaceStreamMediaFinalizeLivestream(ctx context.Context, b return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) } - return &placestream.MediaFinalizeLivestream_Output{UploadId: uploadID}, nil + // Create a draft VOD in the 'processing' state, inheriting the livestream's + // metadata and linking back to it via connections. The client no longer polls + // or publishes — the draft reaches 'ready' server-side and the user + // publishes it later from the Drafts tab. + draft, err := s.createLivestreamDraft(ctx, session.DID, uploadID, ls) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + + return &placestream.MediaFinalizeLivestream_Output{UploadId: uploadID, DraftUri: draft.URI}, nil } diff --git a/pkg/spxrpc/place_stream_vod_drafts.go b/pkg/spxrpc/place_stream_vod_drafts.go index 89ef8198..0b7e5258 100644 --- a/pkg/spxrpc/place_stream_vod_drafts.go +++ b/pkg/spxrpc/place_stream_vod_drafts.go @@ -6,8 +6,11 @@ import ( "net/http" "time" + comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/labstack/echo/v4" "github.com/streamplace/oatproxy/pkg/oatproxy" + "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/statedb" placestream "stream.place/streamplace/pkg/streamplace" "stream.place/streamplace/pkg/vod" ) @@ -188,3 +191,64 @@ func (s *Server) handlePlaceStreamVodPublishDraft(ctx context.Context, body *pla } return &placestream.VodPublishDraft_Output{VideoUri: videoURI, VideoCid: videoCID}, nil } + +// createLivestreamDraft builds a 'processing'-state draft that inherits the +// livestream's title/activity/tags and links back to the livestream record via +// connections. Called by the finalizeLivestream handler at kickoff so the user +// can navigate straight to the draft while processing runs server-side. +func (s *Server) createLivestreamDraft(ctx context.Context, did, uploadID string, ls *model.Livestream) (*statedb.DraftVideo, error) { + view, err := ls.ToLivestreamView() + if err != nil { + return nil, err + } + rec, ok := view.Record.Val.(*placestream.Livestream) + if !ok { + return nil, errors.New("livestream record is not a place.stream.livestream") + } + + title := rec.Title + if title == "" { + title = "Livestream" + } + draftRec := &placestream.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: title, + Status: "processing", + CreatedAt: time.Now().UTC().Format(time.RFC3339), + } + if rec.Activity != nil { + draftRec.Activity = livestreamActivityToDraft(rec.Activity) + } + if len(rec.Tags) > 0 { + draftRec.Tags = rec.Tags + } + // Link back to the source livestream so a published VOD carries the + // connection (the existing UI uses this to flip a finalized row to + // "View VOD"). + draftRec.Connections = []*placestream.VodDraftVideo_Connections_Elem{{ + Video_Connection: &placestream.Video_Connection{ + LexiconTypeID: "place.stream.video#connection", + Ref: &comatproto.RepoStrongRef{ + Uri: ls.URI, + Cid: ls.CID, + }, + }, + }} + + return s.statefulDB.CreateDraft(ctx, did, uploadID, draftRec) +} + +// livestreamActivityToDraft maps a livestream's activity union onto the draft's +// (they reference the same defs). +func livestreamActivityToDraft(a *placestream.Livestream_Activity) *placestream.VodDraftVideo_Activity { + if a == nil { + return nil + } + if a.Defs_ActivityGame != nil { + return &placestream.VodDraftVideo_Activity{Defs_ActivityGame: a.Defs_ActivityGame} + } + if a.Defs_ActivityLabel != nil { + return &placestream.VodDraftVideo_Activity{Defs_ActivityLabel: a.Defs_ActivityLabel} + } + return nil +} diff --git a/pkg/statedb/draft_video.go b/pkg/statedb/draft_video.go index 744e997f..397fd7ec 100644 --- a/pkg/statedb/draft_video.go +++ b/pkg/statedb/draft_video.go @@ -3,10 +3,12 @@ package statedb import ( "bytes" "context" + "encoding/json" "errors" "fmt" "time" + comatproto "github.com/bluesky-social/indigo/api/atproto" "gorm.io/gorm" "stream.place/streamplace/pkg/spid" "stream.place/streamplace/pkg/streamplace" @@ -206,6 +208,54 @@ 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. +func (state *StatefulDB) markDraftReadyFromUpload(ctx context.Context, uploadID string) error { + upload, err := state.GetUpload(ctx, uploadID) + if err != nil { + return fmt.Errorf("get upload: %w", err) + } + if upload == nil { + return nil + } + source, err := draftSourceFromTrackURIs(upload.TrackURIs) + if err != nil { + return err + } + return state.SetDraftReady(ctx, uploadID, source, upload.DurationMS, upload.ContentCID) +} + +// draftSourceFromTrackURIs builds the draft's source union (sourceTracks) from +// the {uri,cid} JSON array the Upload row stores — the same conversion +// vod.sourceTracksFromUpload performs for the publishVideo path, replicated +// here because pkg/statedb can't import pkg/vod (cycle). +func draftSourceFromTrackURIs(trackURIsJSON string) (*streamplace.VodDraftVideo_Source, error) { + if trackURIsJSON == "" { + return nil, nil + } + var refs []struct { + URI string `json:"uri"` + CID string `json:"cid"` + } + if err := json.Unmarshal([]byte(trackURIsJSON), &refs); err != nil { + return nil, fmt.Errorf("decode track refs: %w", err) + } + tracks := make([]*comatproto.RepoStrongRef, 0, len(refs)) + for _, r := range refs { + tracks = append(tracks, &comatproto.RepoStrongRef{Uri: r.URI, Cid: r.CID}) + } + return &streamplace.VodDraftVideo_Source{ + MediaDefs_SourceTracks: &streamplace.MediaDefs_SourceTracks{ + LexiconTypeID: "place.stream.media.defs#sourceTracks", + Tracks: tracks, + }, + }, nil +} + // SetDraftError flips a draft to status "error" with an error message. func (state *StatefulDB) SetDraftError(ctx context.Context, originUploadID, errMsg string) error { return state.updateDraftByUpload(ctx, originUploadID, func(rec *streamplace.VodDraftVideo) { diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index a685beae..2adf648e 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -194,6 +194,11 @@ func (state *StatefulDB) processVODProcessTask(ctx context.Context, task *AppTas 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) } + // Flip any tied draft to 'error' too, so the user sees the failure in + // the Drafts tab instead of an indefinite 'processing' state. + if derr := state.SetDraftError(ctx, t.UploadID, err.Error()); derr != nil { + log.Warn(ctx, "failed to mark draft as failed", "uploadId", t.UploadID, "error", derr) + } // Complete the task so it doesn't retry — most VOD failures are // permanent (unsupported codec, corrupted file, etc.). _ = state.CompleteTask(ctx, task.ID) @@ -203,6 +208,14 @@ func (state *StatefulDB) processVODProcessTask(ctx context.Context, task *AppTas // (e.g. a publish-records track error) can't be tied to an upload. return fmt.Errorf("vod processing upload %s: %w", t.UploadID, err) } + // The processor (vod.ProcessVOD) calls SetUploadProcessed deep inside its + // own publish path, so by the time it returns the Upload row carries the + // finished TrackURIs / DurationMS / ContentCID. Re-read them and flip the + // tied draft to 'ready'. A missing draft (pre-drafts-era upload, or one + // whose create failed) is a no-op, not an error. + if err := state.markDraftReadyFromUpload(ctx, t.UploadID); err != nil { + log.Warn(ctx, "failed to mark draft ready", "uploadId", t.UploadID, "error", err) + } log.Log(ctx, "vod processed", "uploadId", t.UploadID, "cid", cid) return state.CompleteTask(ctx, task.ID) } @@ -237,11 +250,20 @@ func (state *StatefulDB) processFinalizeLivestreamVODTask(ctx context.Context, t 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) } + // Flip the tied draft to 'error' so the user sees the failure. + if derr := state.SetDraftError(ctx, t.UploadID, err.Error()); derr != nil { + log.Warn(ctx, "failed to mark draft as failed", "uploadId", t.UploadID, "error", derr) + } // 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) } + // As with VODProcess, the finalizer calls SetUploadProcessed internally, so + // re-read the finished Upload row and flip the tied draft to 'ready'. + if err := state.markDraftReadyFromUpload(ctx, t.UploadID); err != nil { + log.Warn(ctx, "failed to mark draft ready", "uploadId", t.UploadID, "error", err) + } log.Log(ctx, "livestream VOD finalized", "uploadId", t.UploadID, "cid", cid) return state.CompleteTask(ctx, task.ID) } diff --git a/pkg/statedb/queue_processor_draft_test.go b/pkg/statedb/queue_processor_draft_test.go new file mode 100644 index 00000000..7df6f15d --- /dev/null +++ b/pkg/statedb/queue_processor_draft_test.go @@ -0,0 +1,117 @@ +package statedb + +import ( + "context" + "encoding/json" + "testing" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/streamplace" +) + +// TestDraftLifecycleThroughVODProcessor is the gating test for the draft +// lifecycle wiring: it proves the queue handler re-reads the Upload row (which +// the real vodProcessor populates via SetUploadProcessed internally) after the +// processor returns just a cid, and writes source/durationMs/content_cid onto +// the draft. A fake processor that returns only a cid must still produce a +// 'ready' draft carrying the upload's source/durationMs. +func TestDraftLifecycleThroughVODProcessor(t *testing.T) { + WithAllDatabases(t, func(state *StatefulDB) { + ctx := context.Background() + did := "did:plc:draft" + + // An upload + its processing draft (as createUpload/onComplete would make). + require.NoError(t, state.CreateUpload(ctx, &Upload{ + ID: "up-lifecycle", RepoDID: did, MimeType: "video/mp4", Backend: "file", + })) + _, err := state.CreateDraft(ctx, did, "up-lifecycle", &streamplace.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: "from upload", + Status: "processing", + CreatedAt: "2026-01-01T00:00:00Z", + }) + require.NoError(t, err) + + // 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"}]` + state.SetVODProcessor(func(ctx context.Context, task VODProcessTask) (string, error) { + require.NoError(t, state.SetUploadProcessed(ctx, task.UploadID, trackURIs, 98765, "muxlcid-lifecycle")) + return "muxlcid-lifecycle", nil + }) + + // Drive one task through the processor directly. processVODProcessTask + // also calls CompleteTask at the end, which no-ops on a synthetic task + // ID with no DB row — that's fine; the draft update happens first. + task := &AppTask{ID: 1, Type: TaskVODProcess, Payload: mustMarshal(t, VODProcessTask{ + UploadID: "up-lifecycle", RepoDID: did, MimeType: "video/mp4", Backend: "file", + })} + _ = state.processVODProcessTask(ctx, task) + + // The draft must now be 'ready' and carry the upload's values. + dv, err := state.GetDraftByUpload(ctx, "up-lifecycle") + require.NoError(t, err) + require.NotNil(t, dv) + require.Equal(t, "muxlcid-lifecycle", dv.ContentCID, "content_cid must be backfilled from the upload row") + rec, err := unmarshalDraft(dv.Data) + require.NoError(t, err) + 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) + }) +} + +// TestDraftLifecycleErrorFlipsDraft verifies a failed processor flips the tied +// draft to 'error' so the user isn't stuck looking at 'processing' forever. +func TestDraftLifecycleErrorFlipsDraft(t *testing.T) { + cli := &config.CLI{DBURL: ":memory:"} + mod, err := model.MakeDB(":memory:") + require.NoError(t, err) + state, err := MakeDB(t.Context(), cli, nil, mod) + require.NoError(t, err) + ctx := context.Background() + + require.NoError(t, state.CreateUpload(ctx, &Upload{ID: "up-fail", RepoDID: "did:plc:err", Backend: "file"})) + _, err = state.CreateDraft(ctx, "did:plc:err", "up-fail", &streamplace.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: "will fail", Status: "processing", CreatedAt: "2026-01-01T00:00:00Z", + }) + require.NoError(t, err) + + state.SetVODProcessor(func(ctx context.Context, task VODProcessTask) (string, error) { + return "", errBoom + }) + + task := &AppTask{ID: 2, Type: TaskVODProcess, Payload: mustMarshal(t, VODProcessTask{ + UploadID: "up-fail", RepoDID: "did:plc:err", Backend: "file", + })} + // processVODProcessTask returns the wrapped error (logged upstream); that's fine. + _ = state.processVODProcessTask(ctx, task) + + dv, err := state.GetDraftByUpload(ctx, "up-fail") + require.NoError(t, err) + rec, err := unmarshalDraft(dv.Data) + require.NoError(t, err) + require.Equal(t, "error", rec.Status) + require.NotNil(t, rec.Error) + require.Equal(t, "boom", *rec.Error) +} + +var errBoom = errString("boom") + +type errString string + +func (e errString) Error() string { return string(e) } + +func mustMarshal(t *testing.T, v any) []byte { + t.Helper() + b, err := json.Marshal(v) + require.NoError(t, err) + return b +} diff --git a/pkg/upload/upload.go b/pkg/upload/upload.go index 7ab4f7d5..8b63e202 100644 --- a/pkg/upload/upload.go +++ b/pkg/upload/upload.go @@ -33,6 +33,7 @@ import ( "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/statedb" + "stream.place/streamplace/pkg/streamplace" ) // Backend identifies which TUS data store backs an upload. @@ -325,6 +326,30 @@ func (m *Manager) onComplete(ctx context.Context, ev tushandler.HookEvent) { if _, err := m.state.EnqueueTask(ctx, statedb.TaskVODProcess, payload, statedb.WithTaskKey("vod-process:"+row.ID)); err != nil { log.Error(ctx, "failed to enqueue vod-process task", "error", err) } + + // Create a draft VOD in the 'processing' state. Unlike finalize (which + // inherits livestream metadata), a plain upload starts with placeholder + // metadata the user fills in from the Drafts tab once processing completes. + // A failed create here is non-fatal: the upload still processes, it just + // won't surface as a draft (the user re-publishes via publishVideo). + if _, err := m.state.CreateDraft(ctx, row.RepoDID, row.ID, &streamplace.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: filenameOrDefault(row.Filename), + Status: "processing", + CreatedAt: time.Now().UTC().Format(time.RFC3339), + }); err != nil { + log.Error(ctx, "failed to create draft for upload", "error", err) + } +} + +// filenameOrDefault returns a human placeholder title for a freshly-uploaded +// draft: the original filename if present, else a generic "Uploaded video". +// The user edits this before publishing. +func filenameOrDefault(filename string) string { + if filename != "" { + return filename + } + return "Uploaded video" } func (m *Manager) mintToken(did, uploadID string, expiresAt time.Time) (string, error) { diff --git a/pkg/vod/publish_draft_test.go b/pkg/vod/publish_draft_test.go index 22a32caf..70143685 100644 --- a/pkg/vod/publish_draft_test.go +++ b/pkg/vod/publish_draft_test.go @@ -43,9 +43,9 @@ func TestPublishDraftNotFoundForOtherUsersDraft(t *testing.T) { // Alice owns the draft; Bob tries to publish it. dv, err := state.CreateDraft(ctx, "did:plc:alice", "up-1", &streamplace.VodDraftVideo{ LexiconTypeID: "place.stream.vod.draftVideo", - Title: "Alice's draft", - Status: "ready", - CreatedAt: "2026-01-01T00:00:00Z", + Title: "Alice's draft", + Status: "ready", + CreatedAt: "2026-01-01T00:00:00Z", }) require.NoError(t, err)