diff --git a/pkg/spxrpc/place_stream_vod_drafts.go b/pkg/spxrpc/place_stream_vod_drafts.go new file mode 100644 index 000000000..89ef8198c --- /dev/null +++ b/pkg/spxrpc/place_stream_vod_drafts.go @@ -0,0 +1,190 @@ +package spxrpc + +import ( + "context" + "errors" + "net/http" + "time" + + "github.com/labstack/echo/v4" + "github.com/streamplace/oatproxy/pkg/oatproxy" + placestream "stream.place/streamplace/pkg/streamplace" + "stream.place/streamplace/pkg/vod" +) + +// draft helpers ───────────────────────────────────────────────────────────── + +// draftNotFound and draftForbidden map the common ownership/not-found case to a +// 404, so a caller can't probe another user's draft URIs. +func draftNotFound() error { + return echo.NewHTTPError(http.StatusNotFound, "draft not found") +} + +// loadOwnedDraft fetches a draft by URI and verifies it belongs to session.DID. +// Returns a 404 (draftNotFound) whether the draft is absent or belongs to +// someone else, so a caller can't probe another user's draft URIs. +func (s *Server) loadOwnedDraft(ctx context.Context, uri, did string) (*placestream.VodDraftDefs_DraftView, error) { + if uri == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "uri is required") + } + dv, err := s.statefulDB.GetDraft(ctx, uri) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + if dv == nil || dv.UserDID != did { + return nil, draftNotFound() + } + view, err := dv.ToDraftView() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + return view, nil +} + +// handlers ───────────────────────────────────────────────────────────────── + +func (s *Server) handlePlaceStreamVodListDrafts(ctx context.Context, cursor string, limit int) (*placestream.VodListDrafts_Output, error) { + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required") + } + drafts, err := s.statefulDB.ListDrafts(ctx, session.DID, limit, cursor) + if err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, err.Error()) + } + views := make([]*placestream.VodDraftDefs_DraftView, 0, len(drafts)) + for _, dv := range drafts { + v, err := dv.ToDraftView() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + views = append(views, v) + } + out := &placestream.VodListDrafts_Output{Drafts: views} + if len(drafts) > 0 { + c := drafts[len(drafts)-1].CreatedAt.Format(time.RFC3339Nano) + out.Cursor = &c + } + return out, nil +} + +func (s *Server) handlePlaceStreamVodGetDraft(ctx context.Context, uri string) (*placestream.VodGetDraft_Output, error) { + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required") + } + owned, err := s.loadOwnedDraft(ctx, uri, session.DID) + if err != nil { + return nil, err + } + return &placestream.VodGetDraft_Output{Draft: owned}, nil +} + +func (s *Server) handlePlaceStreamVodUpdateDraft(ctx context.Context, body *placestream.VodUpdateDraft_Input) (*placestream.VodUpdateDraft_Output, error) { + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required") + } + if body.Uri == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "uri is required") + } + // Verify ownership before mutating. + if _, err := s.loadOwnedDraft(ctx, body.Uri, session.DID); err != nil { + return nil, err + } + // Apply only the editable fields present in the partial input. The server- + // authoritative fields (source, durationMs, status, error) are never touched + // here — the closure only mutates editable metadata. + updated, err := s.statefulDB.UpdateDraftMetadata(ctx, body.Uri, func(rec *placestream.VodDraftVideo) { + if body.Title != nil { + rec.Title = *body.Title + } + if body.Description != nil { + rec.Description = body.Description + } + // descriptionFacets / tags: nil means "not provided" (leave as-is); an + // empty slice is a deliberate clear, so always overwrite when non-nil. + if body.DescriptionFacets != nil { + rec.DescriptionFacets = body.DescriptionFacets + } + if body.Tags != nil { + rec.Tags = body.Tags + } + if body.Thumb != nil { + rec.Thumb = body.Thumb + } + if body.Activity != nil { + rec.Activity = (*placestream.VodDraftVideo_Activity)(body.Activity) + } + if body.ContentWarnings != nil { + rec.ContentWarnings = body.ContentWarnings + } + if body.ContentRights != nil { + rec.ContentRights = body.ContentRights + } + }) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + if updated == nil { + return nil, draftNotFound() + } + view, err := updated.ToDraftView() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + return &placestream.VodUpdateDraft_Output{Draft: view}, nil +} + +func (s *Server) handlePlaceStreamVodDeleteDraft(ctx context.Context, body *placestream.VodDeleteDraft_Input) error { + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return echo.NewHTTPError(http.StatusUnauthorized, "oauth session required") + } + if body.Uri == "" { + return echo.NewHTTPError(http.StatusBadRequest, "uri is required") + } + // Verify ownership before deleting. + if _, err := s.loadOwnedDraft(ctx, body.Uri, session.DID); err != nil { + return err + } + deleted, err := s.statefulDB.DeleteDraft(ctx, body.Uri) + if err != nil { + return echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + if !deleted { + return draftNotFound() + } + return nil +} + +func (s *Server) handlePlaceStreamVodPublishDraft(ctx context.Context, body *placestream.VodPublishDraft_Input) (*placestream.VodPublishDraft_Output, error) { + session, _ := oatproxy.GetOAuthSession(ctx) + if session == nil { + return nil, echo.NewHTTPError(http.StatusUnauthorized, "oauth session required") + } + if body.Uri == "" { + return nil, echo.NewHTTPError(http.StatusBadRequest, "uri is required") + } + if s.playbackStore == nil { + return nil, echo.NewHTTPError(http.StatusServiceUnavailable, "playback store not configured") + } + 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") + } + + videoURI, videoCID, err := vod.PublishDraft(ctx, s.statefulDB, s.playbackStore, session.DID, body.Uri) + if err != nil { + switch { + case errors.Is(err, vod.ErrDraftNotFound): + return nil, draftNotFound() + case errors.Is(err, vod.ErrDraftNotReady): + return nil, echo.NewHTTPError(http.StatusConflict, "draft is not ready to publish") + default: + return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } + } + return &placestream.VodPublishDraft_Output{VideoUri: videoURI, VideoCid: videoCID}, nil +} diff --git a/pkg/spxrpc/place_stream_vod_drafts_stubs.go b/pkg/spxrpc/place_stream_vod_drafts_stubs.go deleted file mode 100644 index f778146d0..000000000 --- a/pkg/spxrpc/place_stream_vod_drafts_stubs.go +++ /dev/null @@ -1,36 +0,0 @@ -package spxrpc - -import ( - "context" - "net/http" - - "github.com/labstack/echo/v4" - placestream "stream.place/streamplace/pkg/streamplace" -) - -// errNotImplemented is returned by the draft XRPC stubs until Phase 3 lands. -var errNotImplemented = echo.NewHTTPError(http.StatusNotImplemented, "draft VOD endpoint not yet implemented") - -// Stubs for the place.stream.vod.* draft XRPCs. Full implementations land in -// Phase 3 once the statedb draft storage layer exists; these satisfy the -// generated server stubs (stubs.go) so the tree compiles in the meantime. - -func (s *Server) handlePlaceStreamVodListDrafts(ctx context.Context, cursor string, limit int) (*placestream.VodListDrafts_Output, error) { - return nil, errNotImplemented -} - -func (s *Server) handlePlaceStreamVodGetDraft(ctx context.Context, uri string) (*placestream.VodGetDraft_Output, error) { - return nil, errNotImplemented -} - -func (s *Server) handlePlaceStreamVodUpdateDraft(ctx context.Context, body *placestream.VodUpdateDraft_Input) (*placestream.VodUpdateDraft_Output, error) { - return nil, errNotImplemented -} - -func (s *Server) handlePlaceStreamVodDeleteDraft(ctx context.Context, body *placestream.VodDeleteDraft_Input) error { - return errNotImplemented -} - -func (s *Server) handlePlaceStreamVodPublishDraft(ctx context.Context, body *placestream.VodPublishDraft_Input) (*placestream.VodPublishDraft_Output, error) { - return nil, errNotImplemented -} diff --git a/pkg/statedb/draft_video.go b/pkg/statedb/draft_video.go index e20de254c..744e997f2 100644 --- a/pkg/statedb/draft_video.go +++ b/pkg/statedb/draft_video.go @@ -257,8 +257,8 @@ func (state *StatefulDB) DeleteDraft(ctx context.Context, uri string) (bool, err return res.RowsAffected > 0, nil } -// toDraftView converts a stored row + its CBOR record into the lexicon view. -func (dv *DraftVideo) toDraftView() (*streamplace.VodDraftDefs_DraftView, error) { +// ToDraftView converts a stored row + its CBOR record into the lexicon view. +func (dv *DraftVideo) ToDraftView() (*streamplace.VodDraftDefs_DraftView, error) { rec, err := unmarshalDraft(dv.Data) if err != nil { return nil, err diff --git a/pkg/vod/publish_draft.go b/pkg/vod/publish_draft.go new file mode 100644 index 000000000..8f7f9f947 --- /dev/null +++ b/pkg/vod/publish_draft.go @@ -0,0 +1,194 @@ +package vod + +import ( + "bytes" + "context" + "errors" + "fmt" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/xrpc" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/trace" + + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/constants" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/spid" + "stream.place/streamplace/pkg/statedb" + "stream.place/streamplace/pkg/streamplace" +) + +// ErrDraftNotFound / ErrDraftNotReady let the publishDraft XRPC handler map +// PublishDraft failures onto the right HTTP status (404 / 409) without +// reaching into this package's internals. +var ( + ErrDraftNotFound = errors.New("draft not found") + ErrDraftNotReady = errors.New("draft is not ready to publish") +) + +// PublishDraft promotes a ready draft VOD to a public place.stream.video +// record in the author's repo, then deletes the draft. It mirrors +// PublishVideo's tail: the draft's editable fields (title, description, tags, +// activity, connections, content warnings/rights, thumb) plus its +// server-authoritative fields (source, durationMs) are carried over into the +// video record, createdAt is set server-side, and a thumbnail is backfilled +// from the content blob (via the draft row's ContentCID) when the draft +// carries none. After a successful putRecord the draft row is deleted. +// +// did is the authenticated user; draftURI is the ats:// URI of their draft. +func PublishDraft(ctx context.Context, state *statedb.StatefulDB, store blob.Store, did, draftURI string) (string, string, error) { + ctx, span := vodTracer.Start(ctx, "vod.PublishDraft", trace.WithAttributes( + attribute.String("did", did), + attribute.String("draft_uri", draftURI), + )) + defer span.End() + + dv, err := state.GetDraft(ctx, draftURI) + if err != nil { + span.RecordError(err) + return "", "", fmt.Errorf("get draft: %w", err) + } + if dv == nil || dv.UserDID != did { + return "", "", ErrDraftNotFound + } + rec, err := unmarshalDraftRecord(dv.Data) + if err != nil { + span.RecordError(err) + return "", "", err + } + if rec.Status != "ready" { + span.SetStatus(codes.Error, "not ready") + return "", "", fmt.Errorf("%w: status is %q", ErrDraftNotReady, rec.Status) + } + + // Build the public video record from the draft. The draft's editable + // fields map 1:1; source/durationMs come straight from the CBOR body (the + // server filled them at SetDraftReady time). createdAt is set here, as + // PublishVideo does. + video := &streamplace.Video{ + LexiconTypeID: constants.PLACE_STREAM_VIDEO, + Title: rec.Title, + Description: rec.Description, + DurationMs: derefInt64(rec.DurationMs), + CreatedAt: time.Now().UTC().Format(time.RFC3339), + } + if rec.Description != nil { + video.Description = rec.Description + } + video.DescriptionFacets = rec.DescriptionFacets + video.Tags = rec.Tags + video.Connections = draftConnectionsToVideo(rec.Connections) + video.Activity = draftActivityToVideo(rec.Activity) + video.ContentWarnings = rec.ContentWarnings + 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 { + video.Source = &streamplace.Video_Source{ + MediaDefs_SourceTracks: rec.Source.MediaDefs_SourceTracks, + } + } else if rec.Source != nil && rec.Source.MediaDefs_SourceClip != nil { + video.Source = &streamplace.Video_Source{ + MediaDefs_SourceClip: rec.Source.MediaDefs_SourceClip, + } + } + + client, err := getUserXRPCClient(ctx, state, did) + if err != nil { + span.RecordError(err) + return "", "", fmt.Errorf("get user xrpc client: %w", err) + } + + // Backfill a thumbnail when the draft didn't carry one, using the content + // blob (MUXL CID on the draft row). Non-fatal: publish without it on any + // failure, matching PublishVideo's behavior. + if video.Thumb == nil && dv.ContentCID != "" { + thumb, terr := generateAndUploadThumbnail(ctx, client, store, dv.ContentCID) + if terr != nil { + log.Warn(ctx, "publishDraft: thumbnail backfill failed; publishing without thumb", "error", terr) + } else { + video.Thumb = thumb + } + } + span.SetAttributes(attribute.Bool("thumb_present", video.Thumb != nil)) + + rkey := spid.TIDClock.Next().String() + inp := comatproto.RepoPutRecord_Input{ + Collection: constants.PLACE_STREAM_VIDEO, + Record: &lexutil.LexiconTypeDecoder{Val: video}, + Rkey: rkey, + Repo: did, + } + out := comatproto.RepoPutRecord_Output{} + if err := client.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.putRecord", map[string]any{}, inp, &out); err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, "put_record") + return "", "", fmt.Errorf("putRecord video: %w", err) + } + span.SetAttributes(attribute.String("uri", out.Uri), attribute.String("cid", out.Cid)) + log.Log(ctx, "published video record from draft", + "uri", out.Uri, "draftUri", draftURI, "title", video.Title) + + // Draft fulfilled — delete it now that the public record exists. + if _, err := state.DeleteDraft(ctx, draftURI); err != nil { + // The video was published successfully; a delete failure here leaves a + // stale draft the user can just discard. Log and proceed. + log.Warn(ctx, "publishDraft: failed to delete draft after publish", "draftUri", draftURI, "error", err) + } + + return out.Uri, out.Cid, nil +} + +// unmarshalDraftRecord CBOR-decodes a draft record body. (Sibling of +// statedb.unmarshalDraft, exposed here so PublishDraft stays in pkg/vod.) +func unmarshalDraftRecord(data []byte) (*streamplace.VodDraftVideo, error) { + var rec streamplace.VodDraftVideo + if err := rec.UnmarshalCBOR(bytes.NewReader(data)); err != nil { + return nil, fmt.Errorf("unmarshal draft CBOR: %w", err) + } + return &rec, nil +} + +func derefInt64(p *int64) int64 { + if p == nil { + return 0 + } + return *p +} + +// draftConnectionsToVideo maps the draft's connections union slice to the video +// record's connections slice (both wrap place.stream.video#connection). +func draftConnectionsToVideo(in []*streamplace.VodDraftVideo_Connections_Elem) []*streamplace.Video_Connections_Elem { + if len(in) == 0 { + return nil + } + out := make([]*streamplace.Video_Connections_Elem, 0, len(in)) + for _, c := range in { + if c == nil || c.Video_Connection == nil { + continue + } + out = append(out, &streamplace.Video_Connections_Elem{ + Video_Connection: c.Video_Connection, + }) + } + return out +} + +// draftActivityToVideo maps the draft's activity union to the video record's. +func draftActivityToVideo(in *streamplace.VodDraftVideo_Activity) *streamplace.Video_Activity { + if in == nil { + return nil + } + if in.Defs_ActivityGame != nil { + return &streamplace.Video_Activity{Defs_ActivityGame: in.Defs_ActivityGame} + } + if in.Defs_ActivityLabel != nil { + return &streamplace.Video_Activity{Defs_ActivityLabel: in.Defs_ActivityLabel} + } + return nil +} diff --git a/pkg/vod/publish_draft_test.go b/pkg/vod/publish_draft_test.go new file mode 100644 index 000000000..22a32caf2 --- /dev/null +++ b/pkg/vod/publish_draft_test.go @@ -0,0 +1,109 @@ +package vod + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/blob" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/statedb" + "stream.place/streamplace/pkg/streamplace" +) + +// newDraftTestState builds an in-memory statedb for draft tests that don't need +// gstreamer or a real XRPC/OAuth session. PublishDraft's early error paths +// (ownership, status) return before any of that is touched. +func newDraftTestState(t *testing.T) *statedb.StatefulDB { + t.Helper() + cli := &config.CLI{DBURL: ":memory:"} + mod, err := model.MakeDB(":memory:") + require.NoError(t, err) + state, err := statedb.MakeDB(t.Context(), cli, nil, mod) + require.NoError(t, err) + return state +} + +func TestPublishDraftNotFoundForMissingDraft(t *testing.T) { + state := newDraftTestState(t) + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + + _, _, err = PublishDraft(t.Context(), state, store, "did:plc:alice", "ats://did:plc:alice/place.stream.vod.drafts/self/did:plc:alice/place.stream.vod.draftVideo/none") + require.ErrorIs(t, err, ErrDraftNotFound) +} + +func TestPublishDraftNotFoundForOtherUsersDraft(t *testing.T) { + state := newDraftTestState(t) + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + ctx := context.Background() + + // 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", + }) + require.NoError(t, err) + + _, _, err = PublishDraft(ctx, state, store, "did:plc:bob", dv.URI) + require.ErrorIs(t, err, ErrDraftNotFound, "a foreign caller must not learn the draft exists") +} + +func TestPublishDraftNotReadyWhileProcessing(t *testing.T) { + state := newDraftTestState(t) + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + ctx := context.Background() + + dv, err := state.CreateDraft(ctx, "did:plc:alice", "up-2", &streamplace.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: "Still cooking", + Status: "processing", + CreatedAt: "2026-01-01T00:00:00Z", + }) + require.NoError(t, err) + + _, _, err = PublishDraft(ctx, state, store, "did:plc:alice", dv.URI) + require.ErrorIs(t, err, ErrDraftNotReady) +} + +func TestPublishDraftNotReadyWhenErrored(t *testing.T) { + state := newDraftTestState(t) + store, err := blob.NewFileStore(t.TempDir()) + require.NoError(t, err) + ctx := context.Background() + + dv, err := state.CreateDraft(ctx, "did:plc:alice", "up-3", &streamplace.VodDraftVideo{ + LexiconTypeID: "place.stream.vod.draftVideo", + Title: "Failed", + Status: "error", + CreatedAt: "2026-01-01T00:00:00Z", + }) + require.NoError(t, err) + + _, _, err = PublishDraft(ctx, state, store, "did:plc:alice", dv.URI) + require.ErrorIs(t, err, ErrDraftNotReady) +} + +// draftConnectionsToVideo / draftActivityToVideo are pure mappers; exercise +// them directly since the full PublishDraft path can't run without an XRPC +// client in a unit test. +func TestDraftConnectionAndActivityMapping(t *testing.T) { + require.Nil(t, draftConnectionsToVideo(nil)) + + conn := &streamplace.Video_Connection{LexiconTypeID: "place.stream.video#connection"} + in := []*streamplace.VodDraftVideo_Connections_Elem{{Video_Connection: conn}, nil} + out := draftConnectionsToVideo(in) + require.Len(t, out, 1) + require.Equal(t, conn, out[0].Video_Connection) + + require.Nil(t, draftActivityToVideo(nil)) + game := &streamplace.Defs_ActivityGame{LexiconTypeID: "place.stream.defs#activityGame"} + act := draftActivityToVideo(&streamplace.VodDraftVideo_Activity{Defs_ActivityGame: game}) + require.NotNil(t, act) + require.Equal(t, game, act.Defs_ActivityGame) +}