From b8f49bcf3ca8b5de2ee4586c51cd466d9aca99bd Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 30 Sep 2026 17:52:50 -0700 Subject: [PATCH] auto-vod: make the handoff idempotent and recoverable - Name the automatic upload after the livestream (UUIDv5 of its URI), so a pass interrupted after creating it or queueing the finalize picks it up again instead of starting a second VOD, and skip livestreams that were already finalized by hand. - Schedule from a redelivered ended record too: scheduling is idempotent, and a redelivery is the retry when it failed the first time. - When publishing a finalized livestream VOD fails, leave a ready draft to publish from the Drafts tab (createLivestreamDraft moves to statedb). - Ignore taps on the app toggle while a save is in flight. --- .../settings/privacy-category-settings.tsx | 10 ++- pkg/atproto/sync.go | 9 ++- .../place_stream_media_finalizelivestream.go | 9 ++- pkg/spxrpc/place_stream_vod_drafts.go | 33 ---------- pkg/statedb/auto_publish_vod.go | 26 +++++++- pkg/statedb/auto_publish_vod_test.go | 31 +++++++++ pkg/statedb/livestream_vod.go | 63 +++++++++++++++---- pkg/statedb/queue_processor.go | 28 ++++++--- pkg/statedb/queue_processor_draft_test.go | 46 ++++++++++++++ 9 files changed, 194 insertions(+), 61 deletions(-) diff --git a/js/app/components/settings/privacy-category-settings.tsx b/js/app/components/settings/privacy-category-settings.tsx index b452c8fd7..d0d8cfe70 100644 --- a/js/app/components/settings/privacy-category-settings.tsx +++ b/js/app/components/settings/privacy-category-settings.tsx @@ -7,7 +7,7 @@ import { zero, } from "@streamplace/components"; import { usePDSAgent } from "@streamplace/components/src/streamplace-store/xrpc"; -import { useEffect, useState } from "react"; +import { useEffect, useRef, useState } from "react"; import { useTranslation } from "react-i18next"; import { ScrollView } from "react-native"; import { useStore } from "store"; @@ -22,6 +22,9 @@ function AutoPublishVodsToggle({ host }: { host: string }) { const toast = useToast(); const agent = usePDSAgent(); const [enabled, setEnabled] = useState(null); + // One write at a time: overlapping writes could finish out of order and + // leave the switch showing the opposite of the last tap. + const saving = useRef(false); useEffect(() => { if (!agent) return; @@ -36,7 +39,8 @@ function AutoPublishVodsToggle({ host }: { host: string }) { } const handleChange = async (value: boolean) => { - if (!agent) return; + if (!agent || saving.current) return; + saving.current = true; setEnabled(value); try { const res = await agent.client.call(place.stream.server.putPreferences, { @@ -53,6 +57,8 @@ function AutoPublishVodsToggle({ host }: { host: string }) { : t("auto-publish-vods-update-failed"), { variant: "error" }, ); + } finally { + saving.current = false; } }; diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 6824f065a..a83e14fd0 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -545,7 +545,14 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD err = atsync.Model.CreateLivestream(ctx, ls) if errors.Is(err, model.ErrAlreadyIndexed) { // Re-announcing an unchanged livestream would light the red circle - // up again and re-queue its finalize task. + // up again and re-queue its finalize task. An ended one still + // schedules its VOD, though: scheduling is idempotent, and a + // redelivery is the retry when scheduling failed the first time. + if !isFirstSync && rec.EndedAt != nil && atsync.CLI.StreamIsAllowed(userDID) == nil { + if err := atsync.StatefulDB.ScheduleAutoPublishVOD(ctx, userDID, aturi.String()); err != nil { + return fmt.Errorf("failed to schedule automatic VOD publishing: %w", err) + } + } return nil } if err != nil { diff --git a/pkg/spxrpc/place_stream_media_finalizelivestream.go b/pkg/spxrpc/place_stream_media_finalizelivestream.go index 49d247292..2f77ddeca 100644 --- a/pkg/spxrpc/place_stream_media_finalizelivestream.go +++ b/pkg/spxrpc/place_stream_media_finalizelivestream.go @@ -5,6 +5,7 @@ import ( "net/http" "strings" + "github.com/google/uuid" "github.com/labstack/echo/v4" "stream.place/streamplace/pkg/log" placestream "stream.place/streamplace/pkg/placestream" @@ -56,10 +57,14 @@ func (s *Server) handlePlaceStreamMediaFinalizeLivestream(ctx context.Context, b return nil, echo.NewHTTPError(http.StatusNotFound, "NoRecording: no completed recording objects for these livestreams") } - uploadID, err := s.statefulDB.CreateLivestreamUpload(ctx, streamer, ordered[0]) + uu, err := uuid.NewV7() if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) } + uploadID := uu.String() + if err := s.statefulDB.CreateLivestreamUpload(ctx, uploadID, streamer, ordered[0]); err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } publish := body.Publish != nil && *body.Publish task := statedb.FinalizeLivestreamVODTask{UploadID: uploadID, RepoDID: streamer, LivestreamURI: ordered[0], LivestreamURIs: ordered} @@ -74,7 +79,7 @@ func (s *Server) handlePlaceStreamMediaFinalizeLivestream(ctx context.Context, b } else { // A draft VOD in the 'processing' state; it reaches 'ready' // server-side and the streamer publishes it from the Drafts tab. - draft, err := s.createLivestreamDraft(ctx, streamer, uploadID, video) + draft, err := s.statefulDB.CreateLivestreamDraft(ctx, streamer, uploadID, video) if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) } diff --git a/pkg/spxrpc/place_stream_vod_drafts.go b/pkg/spxrpc/place_stream_vod_drafts.go index ec1b4b417..b92a0050a 100644 --- a/pkg/spxrpc/place_stream_vod_drafts.go +++ b/pkg/spxrpc/place_stream_vod_drafts.go @@ -10,7 +10,6 @@ import ( "github.com/streamplace/oatproxy/pkg/oatproxy" "stream.place/streamplace/pkg/log" placestream "stream.place/streamplace/pkg/placestream" - "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/vod" ) @@ -220,38 +219,6 @@ 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, v *statedb.VideoDraft) (*statedb.DraftVideo, error) { - draftRec := placestream.VodDraftVideo{ - LexiconTypeID: "place.stream.vod.draftVideo", - Title: v.Title, - Description: v.Description, - Status: "processing", - CreatedAt: time.Now().UTC().Format(time.RFC3339), - } - if v.Activity != nil { - draftRec.Activity = &placestream.VodDraftVideo_Activity{ - Defs_ActivityGame: v.Activity.Defs_ActivityGame, - Defs_ActivityLabel: v.Activity.Defs_ActivityLabel, - } - } - if len(v.Tags) > 0 { - draftRec.Tags = v.Tags - } - // Link back to every source livestream so a published VOD carries the - // connections (the existing UI uses this to flip a finalized row to - // "View VOD", and the replay inherits the records' view totals). - for _, c := range v.Connections { - if c.Video_Connection != nil { - draftRec.Connections = append(draftRec.Connections, placestream.VodDraftVideo_Connections_Elem{Video_Connection: c.Video_Connection}) - } - } - 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 { diff --git a/pkg/statedb/auto_publish_vod.go b/pkg/statedb/auto_publish_vod.go index 8af9d6a38..26e2dbfa4 100644 --- a/pkg/statedb/auto_publish_vod.go +++ b/pkg/statedb/auto_publish_vod.go @@ -6,6 +6,7 @@ import ( "fmt" "time" + "github.com/google/uuid" "stream.place/streamplace/pkg/log" placestream "stream.place/streamplace/pkg/placestream" ) @@ -80,6 +81,24 @@ func (state *StatefulDB) processAutoPublishVODTask(ctx context.Context, task *Ap return state.CompleteTask(ctx, task.ID) } + // The upload is named after the livestream, so a pass interrupted after + // creating it (or after queueing the finalize) picks it up again rather + // than starting a second VOD. Any other upload of the livestream is one + // the streamer (or a moderator) finalized by hand. + uploadID := uuid.NewSHA1(uuid.NameSpaceURL, []byte(ls.URI)).String() + uploads, err := state.ListLivestreamUploads(ctx, ls.URI) + if err != nil { + return fmt.Errorf("list livestream uploads: %w", err) + } + created := false + for _, u := range uploads { + if u.ID != uploadID { + log.Log(ctx, "livestream was already finalized into a VOD; not publishing another", "uploadId", u.ID) + return state.CompleteTask(ctx, task.ID) + } + created = true + } + open, err := state.CountOpenS3Segments(ctx, ls.URI) if err != nil { return fmt.Errorf("count open recording objects: %w", err) @@ -110,9 +129,10 @@ func (state *StatefulDB) processAutoPublishVODTask(ctx context.Context, task *Ap if !ok { return fmt.Errorf("record is not a place.stream.livestream: %s", ls.URI) } - uploadID, err := state.CreateLivestreamUpload(ctx, ls.RepoDID, ls.URI) - if err != nil { - return fmt.Errorf("create upload: %w", err) + if !created { + if err := state.CreateLivestreamUpload(ctx, uploadID, ls.RepoDID, ls.URI); err != nil { + return fmt.Errorf("create upload: %w", err) + } } vodTask := FinalizeLivestreamVODTask{ UploadID: uploadID, diff --git a/pkg/statedb/auto_publish_vod_test.go b/pkg/statedb/auto_publish_vod_test.go index 88849e28e..aee9180a5 100644 --- a/pkg/statedb/auto_publish_vod_test.go +++ b/pkg/statedb/auto_publish_vod_test.go @@ -124,6 +124,28 @@ func TestAutoPublishVODQueuesPublishingFinalize(t *testing.T) { }) } +// A pass that dies after handing off (its task is never completed) runs +// again; it must pick up the VOD it started, not start a second one. +func TestAutoPublishVODRetriedHandoff(t *testing.T) { + WithAllDatabases(t, func(state *StatefulDB) { + ctx := context.Background() + did := "did:plc:interrupted" + uri := endedLivestream(t, state, did) + recordObject(t, state, did, uri, "a.m4s", true) + require.NoError(t, state.ScheduleAutoPublishVOD(ctx, did, uri)) + tasks := pendingTasks(t, state, TaskAutoPublishVOD) + require.Len(t, tasks, 1) + + require.NoError(t, state.processAutoPublishVODTask(ctx, &tasks[0])) + require.NoError(t, state.processAutoPublishVODTask(ctx, &tasks[0])) + + require.Len(t, pendingTasks(t, state, TaskFinalizeLivestreamVOD), 1, "one VOD") + uploads, err := state.ListLivestreamUploads(ctx, uri) + require.NoError(t, err) + require.Len(t, uploads, 1, "one upload") + }) +} + func TestAutoPublishVODWaitsForRecording(t *testing.T) { WithAllDatabases(t, func(state *StatefulDB) { did := "did:plc:stillwriting" @@ -180,5 +202,14 @@ func TestAutoPublishVODSkips(t *testing.T) { require.NoError(t, state.processAutoPublishVODTask(ctx, &tasks[0])) require.Empty(t, pendingTasks(t, state, TaskFinalizeLivestreamVOD)) require.Empty(t, pendingTasks(t, state, TaskAutoPublishVOD)) + + // Finalized by hand (the Finalize button, or a moderator) before + // the task ran. + did = "did:plc:byhand" + uri = endedLivestream(t, state, did) + recordObject(t, state, did, uri, "a.m4s", true) + require.NoError(t, state.CreateLivestreamUpload(ctx, "by-hand", did, uri)) + runAutoPublish(t, state, did, uri) + require.Empty(t, pendingTasks(t, state, TaskFinalizeLivestreamVOD)) }) } diff --git a/pkg/statedb/livestream_vod.go b/pkg/statedb/livestream_vod.go index c59d672a4..776e13645 100644 --- a/pkg/statedb/livestream_vod.go +++ b/pkg/statedb/livestream_vod.go @@ -3,8 +3,8 @@ package statedb import ( "context" "strings" + "time" - "github.com/google/uuid" "stream.place/streamplace/pkg/comatproto" "stream.place/streamplace/pkg/model" placestream "stream.place/streamplace/pkg/placestream" @@ -57,24 +57,63 @@ func VideoDraftForLivestreams(items []LivestreamItem, title, description string) return v } +// livestreamUploadBackend marks the synthetic Upload rows of livestream VODs; +// their Location is the first livestream record of the recording. +const livestreamUploadBackend = "live" + // CreateLivestreamUpload creates the synthetic Upload row a livestream VOD is // finalized into, so clients follow it with the getUploadStatus / publishVideo // flow they already have for resumable uploads. firstURI is the first // livestream record of the recording. -func (state *StatefulDB) CreateLivestreamUpload(ctx context.Context, repoDID, firstURI string) (string, error) { - uu, err := uuid.NewV7() - if err != nil { - return "", err - } - uploadID := uu.String() - if err := state.CreateUpload(ctx, &Upload{ +func (state *StatefulDB) CreateLivestreamUpload(ctx context.Context, uploadID, repoDID, firstURI string) error { + return state.CreateUpload(ctx, &Upload{ ID: uploadID, RepoDID: repoDID, MimeType: "video/mp4", - Backend: "live", + Backend: livestreamUploadBackend, Location: firstURI, - }); err != nil { - return "", err + }) +} + +// ListLivestreamUploads lists the VOD uploads whose recording starts with the +// given livestream record. +func (state *StatefulDB) ListLivestreamUploads(ctx context.Context, firstURI string) ([]Upload, error) { + var out []Upload + err := state.DB.WithContext(ctx). + Where("backend = ? AND location = ?", livestreamUploadBackend, firstURI). + Find(&out).Error + return out, err +} + +// CreateLivestreamDraft creates the draft VOD for a livestream upload, in the +// 'processing' state, inheriting the video record's title, description, +// activity, tags and connections back to the livestream records. The draft +// turns 'ready' with the upload; the streamer publishes it from the Drafts +// tab. +func (state *StatefulDB) CreateLivestreamDraft(ctx context.Context, did, uploadID string, v *VideoDraft) (*DraftVideo, error) { + draftRec := placestream.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: v.Title, + Description: v.Description, + Status: "processing", + CreatedAt: time.Now().UTC().Format(time.RFC3339), + } + if v.Activity != nil { + draftRec.Activity = &placestream.VodDraftVideo_Activity{ + Defs_ActivityGame: v.Activity.Defs_ActivityGame, + Defs_ActivityLabel: v.Activity.Defs_ActivityLabel, + } + } + if len(v.Tags) > 0 { + draftRec.Tags = v.Tags + } + // Link back to every source livestream so a published VOD carries the + // connections (the existing UI uses this to flip a finalized row to + // "View VOD", and the replay inherits the records' view totals). + for _, c := range v.Connections { + if c.Video_Connection != nil { + draftRec.Connections = append(draftRec.Connections, placestream.VodDraftVideo_Connections_Elem{Video_Connection: c.Video_Connection}) + } } - return uploadID, nil + return state.CreateDraft(ctx, did, uploadID, &draftRec) } diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index 03b873666..892e25513 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -109,8 +109,9 @@ type FinalizeLivestreamVODTask struct { // Publish, when set, describes the place.stream.video record to // publish in the streamer's repo (with their stored session) as soon // as the VOD is finalized, instead of leaving a draft for them to - // publish from the app. Set by the operator's finalize route and by - // automatic VOD publishing (AutoPublishVODTask). + // publish from the app (unless publishing fails). Set by the + // operator's finalize route and by automatic VOD publishing + // (AutoPublishVODTask). Publish *VideoDraft `json:"publish,omitempty"` } @@ -374,16 +375,27 @@ func (state *StatefulDB) processFinalizeLivestreamVODTask(ctx context.Context, t } log.Log(ctx, "livestream VOD finalized", "uploadId", t.UploadID, "cid", cid) if t.Publish != nil { - // The VOD is finalized either way: a failed publish leaves an upload - // the streamer can still publish from the app, so it is logged, not - // retried (a retry would finalize and publish the tracks again). + // The VOD is finalized either way, so a failed publish is not + // retried (a retry would finalize the recording again). It leaves + // the streamer a ready draft instead, to publish from the Drafts tab; + // the publish remembers any track records it got as far as minting on + // the upload, and the draft's publish reuses them. + var perr error if state.videoPublisher == nil { - log.Error(ctx, "finalize-livestream-vod: publish requested but no video publisher configured", "uploadId", t.UploadID) - } else if uri, vcid, perr := state.videoPublisher(ctx, t); perr != nil { - log.Error(ctx, "finalize-livestream-vod: VOD finalized but publishing the video record failed", "uploadId", t.UploadID, "error", perr) + perr = errors.New("no video publisher configured") + } else if uri, vcid, err := state.videoPublisher(ctx, t); err != nil { + perr = err } else { log.Log(ctx, "livestream VOD published", "uploadId", t.UploadID, "uri", uri, "cid", vcid) } + if perr != nil { + log.Error(ctx, "finalize-livestream-vod: VOD finalized but publishing the video record failed; leaving a draft", "uploadId", t.UploadID, "error", perr) + if _, err := state.CreateLivestreamDraft(ctx, t.RepoDID, t.UploadID, t.Publish); err != nil { + log.Error(ctx, "finalize-livestream-vod: could not create a draft for the unpublished VOD", "uploadId", t.UploadID, "error", err) + } else if err := state.markDraftReadyFromUpload(ctx, t.UploadID); err != nil { + log.Warn(ctx, "failed to mark draft ready", "uploadId", t.UploadID, "error", err) + } + } } return state.CompleteTask(ctx, task.ID) } diff --git a/pkg/statedb/queue_processor_draft_test.go b/pkg/statedb/queue_processor_draft_test.go index 21f4201c1..836978dd4 100644 --- a/pkg/statedb/queue_processor_draft_test.go +++ b/pkg/statedb/queue_processor_draft_test.go @@ -109,6 +109,52 @@ func TestDraftLifecycleErrorFlipsDraft(t *testing.T) { require.Equal(t, "boom", *rec.Error) } +// A livestream VOD finalized for publishing whose publish fails is left as a +// ready draft, so the streamer can still publish it from the Drafts tab; +// one that publishes leaves no draft. +func TestFinalizeLivestreamVODPublishFailureLeavesDraft(t *testing.T) { + WithAllDatabases(t, func(state *StatefulDB) { + ctx := context.Background() + did := "did:plc:publishfail" + state.SetLivestreamVODFinalizer(func(ctx context.Context, task FinalizeLivestreamVODTask) (string, error) { + require.NoError(t, state.SetUploadProcessed(ctx, task.UploadID, 5000, "muxlcid-"+task.UploadID, "did:key:signing", `{"durationMs":5000}`, 100)) + return "muxlcid-" + task.UploadID, nil + }) + publishErr := error(errBoom) + state.SetVideoPublisher(func(ctx context.Context, task FinalizeLivestreamVODTask) (string, string, error) { + if publishErr != nil { + return "", "", publishErr + } + return "at://" + did + "/place.stream.video/v", "bafyvideo", nil + }) + finalize := func(uploadID string) { + require.NoError(t, state.CreateLivestreamUpload(ctx, uploadID, did, "at://"+did+"/place.stream.livestream/"+uploadID)) + task, err := state.EnqueueTask(ctx, TaskFinalizeLivestreamVOD, FinalizeLivestreamVODTask{ + UploadID: uploadID, RepoDID: did, + Publish: &VideoDraft{Title: "Late night stream", Tags: []string{"chill"}}, + }) + require.NoError(t, err) + require.NoError(t, state.processFinalizeLivestreamVODTask(ctx, task)) + } + + finalize("up-unpublished") + dv, err := state.GetDraftByUpload(ctx, "up-unpublished") + require.NoError(t, err) + require.NotNil(t, dv, "a draft to publish by hand") + rec, err := unmarshalDraft(dv.Data) + require.NoError(t, err) + require.Equal(t, "ready", rec.Status) + require.Equal(t, "Late night stream", rec.Title) + require.Equal(t, []string{"chill"}, rec.Tags) + + publishErr = nil + finalize("up-published") + dv, err = state.GetDraftByUpload(ctx, "up-published") + require.NoError(t, err) + require.Nil(t, dv, "nothing left to publish") + }) +} + var errBoom = errString("boom") type errString string -- 2.51.2