From affb4707707db19a3665a99cd96f2e1128f2095e Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 16 May 2026 14:55:52 -0700 Subject: [PATCH] model: take + return atproto records on Video/MediaTrack/MediaOrigin/BetaInvite MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Mirrors the ModerationDelegation/BroadcastOrigin pattern. Upsert methods now take a *streamplace.X plus the AT-URI and derive URI / CID / RepoDID from those; Get-by-URI returns the decoded record so callers don't have to round-trip through CBOR themselves. Per-list lookups still return model rows because callers need identity context the lexicon record doesn't carry. Drops Title/DurationMS off Video, TrackID/MediaType/Language off MediaTrack, Size/MimeType off MediaOrigin, RKey off all four. The only columns left are the ones we actually index by — every other field is reachable through ToRecord(). As a bonus, unindexed attributes no longer tempt callers into accidental full-table scans. Co-Authored-By: Claude Opus 4.7 --- pkg/atproto/sync.go | 64 +++----------- pkg/model/beta_invite.go | 41 +++++++-- pkg/model/media_origin.go | 69 +++++++++++---- pkg/model/media_track.go | 86 ++++++++++++++----- pkg/model/model.go | 14 +-- pkg/model/video.go | 78 ++++++++++++----- pkg/spxrpc/beta_invite_test.go | 22 ++--- pkg/spxrpc/labels_test.go | 22 +++-- pkg/spxrpc/place_stream_playback_getvideo.go | 27 +++--- .../place_stream_playback_getvideo_test.go | 65 ++++---------- 10 files changed, 274 insertions(+), 214 deletions(-) diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 7ba19bf5..03d39688 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -731,64 +731,30 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD case *streamplace.Video: // Index the video record so playback can resolve it by URI // without needing to round-trip back to the user's PDS. - v := &model.Video{ - URI: aturi.String(), - CID: cid, - RepoDID: userDID, - RKey: rkey.String(), - Title: rec.Title, - Record: *recCBOR, - IndexedAt: now, - } - v.DurationMS = &rec.Duration - if err := atsync.Model.UpsertVideo(ctx, v); err != nil { + if err := atsync.Model.UpsertVideo(ctx, rec, aturi); err != nil { return fmt.Errorf("failed to upsert video: %w", err) } log.Debug(ctx, "indexed video", "uri", aturi.String(), "title", rec.Title) case *streamplace.MediaTrack: - // Pull the muxlTrack subobject — that's where the blob CID - // and track-within-container metadata live. Tracks not - // backed by a muxlTrack (we don't define any other shape - // yet) are skipped with a warning. + // Tracks not backed by a muxlTrack (we don't define any other + // shape yet) are skipped with a warning — there'd be no blob + // to key the row off of. if rec.Track == nil || rec.Track.MediaDefs_MuxlTrack == nil { log.Warn(ctx, "track record missing muxlTrack; skipping", "uri", aturi.String()) return nil } - mt := rec.Track.MediaDefs_MuxlTrack - t := &model.MediaTrack{ - URI: aturi.String(), - CID: cid, - RepoDID: userDID, - RKey: rkey.String(), - Blob: mt.Blob, - TrackID: mt.TrackId, - MediaType: mt.MediaType, - Language: mt.Language, - Record: *recCBOR, - IndexedAt: now, - } - if err := atsync.Model.UpsertMediaTrack(ctx, t); err != nil { + if err := atsync.Model.UpsertMediaTrack(ctx, rec, aturi); err != nil { return fmt.Errorf("failed to upsert media track: %w", err) } + mt := rec.Track.MediaDefs_MuxlTrack log.Debug(ctx, "indexed media track", "uri", aturi.String(), "blob", mt.Blob, "mediaType", mt.MediaType) case *streamplace.MediaOrigin: // Origin records are published by streamplace nodes (not users) - // against their own server-repo DID. The userDID here is the - // publishing server's DID — same field, different semantics. - o := &model.MediaOrigin{ - URI: aturi.String(), - CID: cid, - ServerDID: userDID, - RKey: rkey.String(), - Blob: rec.Blob, - Size: rec.Size, - MimeType: rec.MimeType, - Record: *recCBOR, - IndexedAt: now, - } - if err := atsync.Model.UpsertMediaOrigin(ctx, o); err != nil { + // against their own server-repo DID. The aturi's authority is + // the publishing server. + if err := atsync.Model.UpsertMediaOrigin(ctx, rec, aturi); err != nil { return fmt.Errorf("failed to upsert media origin: %w", err) } log.Debug(ctx, "indexed media origin", "uri", aturi.String(), "blob", rec.Blob, "server", userDID) @@ -799,17 +765,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD // callers filter by RepoDID to a single operator-configured // issuer (the `--beta-invite-did` flag), so anyone else // minting these records is harmless noise. - inv := &model.BetaInvite{ - URI: aturi.String(), - CID: cid, - RepoDID: userDID, - RKey: rkey.String(), - DID: rec.Did, - Feature: rec.Feature, - Record: *recCBOR, - IndexedAt: now, - } - if err := atsync.Model.UpsertBetaInvite(ctx, inv); err != nil { + if err := atsync.Model.UpsertBetaInvite(ctx, rec, aturi); err != nil { return fmt.Errorf("failed to upsert beta invite: %w", err) } log.Debug(ctx, "indexed beta invite", "uri", aturi.String(), "did", rec.Did, "feature", rec.Feature) diff --git a/pkg/model/beta_invite.go b/pkg/model/beta_invite.go index 40abbbb0..329cfc09 100644 --- a/pkg/model/beta_invite.go +++ b/pkg/model/beta_invite.go @@ -1,16 +1,21 @@ package model import ( + "bytes" "context" "fmt" "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/spid" + "stream.place/streamplace/pkg/streamplace" ) -// BetaInvite is the indexed view of a place.stream.beta.invite record. -// One row per (RepoDID, DID, Feature) triple: a single account can -// grant a single feature to a given DID at most once at a time. -// -// The trust model is "we believe whoever owns the record's repo": at +// BetaInvite is the indexed view of a place.stream.beta.invite +// record. One row per (RepoDID, DID, Feature) triple; all three sit +// in the composite index since HasBetaInvite filters by exactly that +// shape. The trust model is "we believe whoever owns the repo": at // the gate, callers must filter by RepoDID to the operator-configured // `--beta-invite-did` so a random user's repo can't mint invites for // our node. @@ -18,15 +23,35 @@ type BetaInvite struct { URI string `gorm:"primaryKey;column:uri"` CID string `gorm:"column:cid"` RepoDID string `gorm:"column:repo_did;index:idx_invites_lookup,priority:1"` - RKey string `gorm:"column:rkey"` DID string `gorm:"column:did;index:idx_invites_lookup,priority:2"` Feature string `gorm:"column:feature;index:idx_invites_lookup,priority:3"` Record []byte `gorm:"column:record"` IndexedAt time.Time `gorm:"column:indexed_at"` } -func (m *DBModel) UpsertBetaInvite(ctx context.Context, v *BetaInvite) error { - return m.DB.WithContext(ctx).Save(v).Error +func (m *DBModel) UpsertBetaInvite(ctx context.Context, rec *streamplace.BetaInvite, aturi syntax.ATURI) error { + repoDID, err := aturi.Authority().AsDID() + if err != nil { + return fmt.Errorf("invalid ATURI authority: %w", err) + } + cid, err := spid.GetCID(rec) + if err != nil { + return fmt.Errorf("get beta invite CID: %w", err) + } + var buf bytes.Buffer + if err := rec.MarshalCBOR(&buf); err != nil { + return fmt.Errorf("marshal beta invite record: %w", err) + } + inv := &BetaInvite{ + URI: aturi.String(), + CID: cid.String(), + RepoDID: repoDID.String(), + DID: rec.Did, + Feature: rec.Feature, + Record: buf.Bytes(), + IndexedAt: aqtime.FromTime(time.Now().UTC()).Time().UTC(), + } + return m.DB.WithContext(ctx).Save(inv).Error } func (m *DBModel) DeleteBetaInvite(ctx context.Context, uri string) error { diff --git a/pkg/model/media_origin.go b/pkg/model/media_origin.go index 48c227ce..410811eb 100644 --- a/pkg/model/media_origin.go +++ b/pkg/model/media_origin.go @@ -1,37 +1,69 @@ package model import ( + "bytes" "context" "errors" "fmt" "time" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" "gorm.io/gorm" + "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/spid" + "stream.place/streamplace/pkg/streamplace" ) // MediaOrigin is the indexed view of a place.stream.media.origin -// record. Each row is an attestation by ServerDID (a streamplace node's -// DID) that it has the blob at the given CID and can serve it via -// XRPC. Many origin rows can point at the same Blob; the playback -// path uses them as a candidate set of nodes to fetch from. -// -// The rkey of the underlying record is conventionally the blob CID -// itself (origin records use key=any). We also store the rkey -// explicitly so the (server, blob) tuple has a stable foreign key -// shape even if that convention changes. +// record: a server's attestation that it holds the blob at the given +// CID. Many origin rows can point at the same Blob (one per server); +// the playback path queries by (Blob, ServerDID) to assemble the +// candidate set of nodes to fetch from. Size/MimeType and any other +// blob metadata stay in the CBOR Record blob. type MediaOrigin struct { URI string `gorm:"primaryKey;column:uri"` CID string `gorm:"column:cid"` ServerDID string `gorm:"column:server_did;index:idx_origins_blob_server,priority:2"` - RKey string `gorm:"column:rkey"` Blob string `gorm:"column:blob;index:idx_origins_blob_server,priority:1"` - Size int64 `gorm:"column:size"` - MimeType string `gorm:"column:mime_type"` Record []byte `gorm:"column:record"` IndexedAt time.Time `gorm:"column:indexed_at"` } -func (m *DBModel) UpsertMediaOrigin(ctx context.Context, o *MediaOrigin) error { +// ToRecord decodes the stored CBOR into the typed lexicon struct. +func (o *MediaOrigin) ToRecord() (*streamplace.MediaOrigin, error) { + rec, err := lexutil.CborDecodeValue(o.Record) + if err != nil { + return nil, fmt.Errorf("decode media origin record: %w", err) + } + origin, ok := rec.(*streamplace.MediaOrigin) + if !ok { + return nil, fmt.Errorf("media origin record decoded as %T, expected *streamplace.MediaOrigin", rec) + } + return origin, nil +} + +func (m *DBModel) UpsertMediaOrigin(ctx context.Context, rec *streamplace.MediaOrigin, aturi syntax.ATURI) error { + serverDID, err := aturi.Authority().AsDID() + if err != nil { + return fmt.Errorf("invalid ATURI authority: %w", err) + } + cid, err := spid.GetCID(rec) + if err != nil { + return fmt.Errorf("get media origin CID: %w", err) + } + var buf bytes.Buffer + if err := rec.MarshalCBOR(&buf); err != nil { + return fmt.Errorf("marshal media origin record: %w", err) + } + o := &MediaOrigin{ + URI: aturi.String(), + CID: cid.String(), + ServerDID: serverDID.String(), + Blob: rec.Blob, + Record: buf.Bytes(), + IndexedAt: aqtime.FromTime(time.Now().UTC()).Time().UTC(), + } return m.DB.WithContext(ctx).Save(o).Error } @@ -39,7 +71,7 @@ func (m *DBModel) DeleteMediaOrigin(ctx context.Context, uri string) error { return m.DB.WithContext(ctx).Where("uri = ?", uri).Delete(&MediaOrigin{}).Error } -func (m *DBModel) GetMediaOriginByURI(ctx context.Context, uri string) (*MediaOrigin, error) { +func (m *DBModel) GetMediaOriginByURI(ctx context.Context, uri string) (*streamplace.MediaOrigin, error) { var o MediaOrigin err := m.DB.WithContext(ctx).Where("uri = ?", uri).First(&o).Error if errors.Is(err, gorm.ErrRecordNotFound) { @@ -48,13 +80,12 @@ func (m *DBModel) GetMediaOriginByURI(ctx context.Context, uri string) (*MediaOr if err != nil { return nil, fmt.Errorf("get media origin by uri: %w", err) } - return &o, nil + return o.ToRecord() } -// GetMediaOriginsByBlob returns every server that has attested to -// hosting the given blob CID. Used by playback to discover candidate -// download sources. Returned in IndexedAt order (newest first) on -// the theory that recent attestations are more reliable. +// GetMediaOriginsByBlob returns every server attestation for the +// given blob CID, newest first. Returns model rows so the caller +// has ServerDID for source selection without a CBOR decode per row. func (m *DBModel) GetMediaOriginsByBlob(ctx context.Context, blob string) ([]*MediaOrigin, error) { var out []*MediaOrigin err := m.DB.WithContext(ctx). diff --git a/pkg/model/media_track.go b/pkg/model/media_track.go index 1f0e029f..0f7596d8 100644 --- a/pkg/model/media_track.go +++ b/pkg/model/media_track.go @@ -1,37 +1,79 @@ package model import ( + "bytes" "context" "errors" "fmt" "time" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" "gorm.io/gorm" + "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/spid" + "stream.place/streamplace/pkg/streamplace" ) -// MediaTrack is the indexed view of a place.stream.media.track record. -// One record per A/V track within a MUXL container; multiple tracks -// share a Blob (the BDASL CID of the container they live in). -// -// Blob is indexed because "list all tracks for this blob" is the -// canonical playback query — having an Origin pointing at a blob -// tells you it's available; having Tracks for that blob tells you -// what's actually in it. +// MediaTrack is the indexed view of a place.stream.media.track +// record. Only fields used to query — Blob (for "list tracks in this +// container") and RepoDID (for ownership checks on those query +// results) — sit beside the URI/CID identity. Everything else lives +// in the CBOR Record blob and is reached via ToRecord. type MediaTrack struct { URI string `gorm:"primaryKey;column:uri"` CID string `gorm:"column:cid"` - RepoDID string `gorm:"column:repo_did;index"` - Repo *Repo `gorm:"foreignKey:DID;references:RepoDID"` - RKey string `gorm:"column:rkey"` + RepoDID string `gorm:"column:repo_did"` Blob string `gorm:"column:blob;index"` - TrackID string `gorm:"column:track_id"` - MediaType string `gorm:"column:media_type;index"` - Language *string `gorm:"column:language"` Record []byte `gorm:"column:record"` IndexedAt time.Time `gorm:"column:indexed_at"` } -func (m *DBModel) UpsertMediaTrack(ctx context.Context, t *MediaTrack) error { +// ToRecord decodes the stored CBOR into the typed lexicon struct. +func (t *MediaTrack) ToRecord() (*streamplace.MediaTrack, error) { + rec, err := lexutil.CborDecodeValue(t.Record) + if err != nil { + return nil, fmt.Errorf("decode media track record: %w", err) + } + track, ok := rec.(*streamplace.MediaTrack) + if !ok { + return nil, fmt.Errorf("media track record decoded as %T, expected *streamplace.MediaTrack", rec) + } + return track, nil +} + +// trackBlob pulls the MUXL container's blob CID off a typed track +// record. Tracks not backed by a muxlTrack (no other shape defined +// yet) return an empty string, leaving it to the caller to decide +// whether that's worth indexing. +func trackBlob(rec *streamplace.MediaTrack) string { + if rec == nil || rec.Track == nil || rec.Track.MediaDefs_MuxlTrack == nil { + return "" + } + return rec.Track.MediaDefs_MuxlTrack.Blob +} + +func (m *DBModel) UpsertMediaTrack(ctx context.Context, rec *streamplace.MediaTrack, aturi syntax.ATURI) error { + repoDID, err := aturi.Authority().AsDID() + if err != nil { + return fmt.Errorf("invalid ATURI authority: %w", err) + } + cid, err := spid.GetCID(rec) + if err != nil { + return fmt.Errorf("get media track CID: %w", err) + } + var buf bytes.Buffer + if err := rec.MarshalCBOR(&buf); err != nil { + return fmt.Errorf("marshal media track record: %w", err) + } + t := &MediaTrack{ + URI: aturi.String(), + CID: cid.String(), + RepoDID: repoDID.String(), + Blob: trackBlob(rec), + Record: buf.Bytes(), + IndexedAt: aqtime.FromTime(time.Now().UTC()).Time().UTC(), + } return m.DB.WithContext(ctx).Save(t).Error } @@ -39,25 +81,25 @@ func (m *DBModel) DeleteMediaTrack(ctx context.Context, uri string) error { return m.DB.WithContext(ctx).Where("uri = ?", uri).Delete(&MediaTrack{}).Error } -func (m *DBModel) GetMediaTrackByURI(ctx context.Context, uri string) (*MediaTrack, error) { +func (m *DBModel) GetMediaTrackByURI(ctx context.Context, uri string) (*streamplace.MediaTrack, error) { var t MediaTrack - err := m.DB.WithContext(ctx).Preload("Repo").Where("uri = ?", uri).First(&t).Error + err := m.DB.WithContext(ctx).Where("uri = ?", uri).First(&t).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, fmt.Errorf("get media track by uri: %w", err) } - return &t, nil + return t.ToRecord() } -// GetMediaTracksByBlob returns every track record that claims to live -// inside the given MUXL container, across all repos. Playback uses -// this to assemble the per-track manifest set for an HLS playlist. +// GetMediaTracksByBlob returns every track row that claims to live +// inside the given MUXL container, across all repos. Returns model +// rows (not decoded records) so callers have RepoDID for ownership +// checks without paying for a CBOR decode per track. func (m *DBModel) GetMediaTracksByBlob(ctx context.Context, blob string) ([]*MediaTrack, error) { var out []*MediaTrack err := m.DB.WithContext(ctx). - Preload("Repo"). Where("blob = ?", blob). Find(&out).Error if err != nil { diff --git a/pkg/model/model.go b/pkg/model/model.go index 61c1dec7..d862d94a 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -133,22 +133,22 @@ type Model interface { GetBadgeIssuanceByURI(ctx context.Context, uri string) (*BadgeIssuance, error) GetBadgeIssuancesForRecipient(ctx context.Context, recipientDID string) ([]*BadgeIssuance, error) - UpsertVideo(ctx context.Context, v *Video) error + UpsertVideo(ctx context.Context, rec *streamplace.Video, aturi syntax.ATURI) error DeleteVideo(ctx context.Context, uri string) error - GetVideoByURI(ctx context.Context, uri string) (*Video, error) + GetVideoByURI(ctx context.Context, uri string) (*streamplace.Video, error) GetLatestVideosForRepo(ctx context.Context, repoDID string, limit int) ([]*Video, error) - UpsertMediaTrack(ctx context.Context, t *MediaTrack) error + UpsertMediaTrack(ctx context.Context, rec *streamplace.MediaTrack, aturi syntax.ATURI) error DeleteMediaTrack(ctx context.Context, uri string) error - GetMediaTrackByURI(ctx context.Context, uri string) (*MediaTrack, error) + GetMediaTrackByURI(ctx context.Context, uri string) (*streamplace.MediaTrack, error) GetMediaTracksByBlob(ctx context.Context, blob string) ([]*MediaTrack, error) - UpsertMediaOrigin(ctx context.Context, o *MediaOrigin) error + UpsertMediaOrigin(ctx context.Context, rec *streamplace.MediaOrigin, aturi syntax.ATURI) error DeleteMediaOrigin(ctx context.Context, uri string) error - GetMediaOriginByURI(ctx context.Context, uri string) (*MediaOrigin, error) + GetMediaOriginByURI(ctx context.Context, uri string) (*streamplace.MediaOrigin, error) GetMediaOriginsByBlob(ctx context.Context, blob string) ([]*MediaOrigin, error) - UpsertBetaInvite(ctx context.Context, v *BetaInvite) error + UpsertBetaInvite(ctx context.Context, rec *streamplace.BetaInvite, aturi syntax.ATURI) error DeleteBetaInvite(ctx context.Context, uri string) error HasBetaInvite(ctx context.Context, fromRepoDID, subjectDID, feature string) (bool, error) } diff --git a/pkg/model/video.go b/pkg/model/video.go index 3c7678f9..266b01f2 100644 --- a/pkg/model/video.go +++ b/pkg/model/video.go @@ -1,34 +1,66 @@ package model import ( + "bytes" "context" "errors" "fmt" "time" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" "gorm.io/gorm" + "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/spid" + "stream.place/streamplace/pkg/streamplace" ) -// Video is the indexed view of a place.stream.video record. Indexing -// happens via the firehose for any user repo we're observing. -// -// The Record column holds the CBOR-encoded record bytes; everything -// else is a denormalized projection used to make playback queries -// (by author, recency, etc.) cheap. Title and Duration are surfaced -// since they're tiny and ubiquitously needed by listing UIs. +// Video is the indexed view of a place.stream.video record. Every +// non-identity field — title, duration, source, etc — lives inside +// the CBOR Record blob; callers decode it via ToRecord (or rely on +// the getters that return *streamplace.Video directly). Only +// indexed fields and the URI/CID identity earn their own column. type Video struct { - URI string `gorm:"primaryKey;column:uri"` - CID string `gorm:"column:cid"` - RepoDID string `gorm:"column:repo_did;index:idx_videos_repo_created,priority:1"` - Repo *Repo `gorm:"foreignKey:DID;references:RepoDID"` - RKey string `gorm:"column:rkey"` - Title string `gorm:"column:title"` - DurationMS *int64 `gorm:"column:duration_ms"` - Record []byte `gorm:"column:record"` - IndexedAt time.Time `gorm:"column:indexed_at;index:idx_videos_repo_created,priority:2"` + URI string `gorm:"primaryKey;column:uri"` + CID string `gorm:"column:cid"` + RepoDID string `gorm:"column:repo_did;index:idx_videos_repo_indexed,priority:1"` + Record []byte `gorm:"column:record"` + IndexedAt time.Time `gorm:"column:indexed_at;index:idx_videos_repo_indexed,priority:2"` } -func (m *DBModel) UpsertVideo(ctx context.Context, v *Video) error { +// ToRecord decodes the stored CBOR into the typed lexicon struct. +func (v *Video) ToRecord() (*streamplace.Video, error) { + rec, err := lexutil.CborDecodeValue(v.Record) + if err != nil { + return nil, fmt.Errorf("decode video record: %w", err) + } + video, ok := rec.(*streamplace.Video) + if !ok { + return nil, fmt.Errorf("video record decoded as %T, expected *streamplace.Video", rec) + } + return video, nil +} + +func (m *DBModel) UpsertVideo(ctx context.Context, rec *streamplace.Video, aturi syntax.ATURI) error { + repoDID, err := aturi.Authority().AsDID() + if err != nil { + return fmt.Errorf("invalid ATURI authority: %w", err) + } + cid, err := spid.GetCID(rec) + if err != nil { + return fmt.Errorf("get video CID: %w", err) + } + var buf bytes.Buffer + if err := rec.MarshalCBOR(&buf); err != nil { + return fmt.Errorf("marshal video record: %w", err) + } + v := &Video{ + URI: aturi.String(), + CID: cid.String(), + RepoDID: repoDID.String(), + Record: buf.Bytes(), + IndexedAt: aqtime.FromTime(time.Now().UTC()).Time().UTC(), + } return m.DB.WithContext(ctx).Save(v).Error } @@ -36,27 +68,27 @@ func (m *DBModel) DeleteVideo(ctx context.Context, uri string) error { return m.DB.WithContext(ctx).Where("uri = ?", uri).Delete(&Video{}).Error } -func (m *DBModel) GetVideoByURI(ctx context.Context, uri string) (*Video, error) { +func (m *DBModel) GetVideoByURI(ctx context.Context, uri string) (*streamplace.Video, error) { var v Video - err := m.DB.WithContext(ctx).Preload("Repo").Where("uri = ?", uri).First(&v).Error + err := m.DB.WithContext(ctx).Where("uri = ?", uri).First(&v).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, fmt.Errorf("get video by uri: %w", err) } - return &v, nil + return v.ToRecord() } -// GetLatestVideosForRepo returns the most recent N videos published -// by the given user. Used to populate per-streamer VOD feeds. +// GetLatestVideosForRepo returns the most recent N video rows by a +// given repo. Model rows (not decoded records) so the caller has +// URI/CID identity for each entry; call ToRecord on a row to decode. func (m *DBModel) GetLatestVideosForRepo(ctx context.Context, repoDID string, limit int) ([]*Video, error) { if limit <= 0 { limit = 25 } var out []*Video err := m.DB.WithContext(ctx). - Preload("Repo"). Where("repo_did = ?", repoDID). Order("indexed_at DESC"). Limit(limit). diff --git a/pkg/spxrpc/beta_invite_test.go b/pkg/spxrpc/beta_invite_test.go index 12a34ace..e422675a 100644 --- a/pkg/spxrpc/beta_invite_test.go +++ b/pkg/spxrpc/beta_invite_test.go @@ -3,12 +3,13 @@ package spxrpc import ( "context" "testing" - "time" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/streamplace" ) const ( @@ -19,16 +20,15 @@ const ( func putInvite(t *testing.T, m model.Model, repoDID, subjectDID, feature string) { t.Helper() - require.NoError(t, m.UpsertBetaInvite(context.Background(), &model.BetaInvite{ - URI: "at://" + repoDID + "/place.stream.beta.invite/" + feature + "-" + subjectDID, - CID: "bafycidplaceholder", - RepoDID: repoDID, - RKey: feature + "-" + subjectDID, - DID: subjectDID, - Feature: feature, - Record: []byte{0xa0}, // empty CBOR map, never decoded in this path - IndexedAt: time.Now().UTC(), - })) + rkey := feature + "-" + subjectDID + aturi, err := syntax.ParseATURI("at://" + repoDID + "/place.stream.beta.invite/" + rkey) + require.NoError(t, err) + require.NoError(t, m.UpsertBetaInvite(context.Background(), &streamplace.BetaInvite{ + LexiconTypeID: "place.stream.beta.invite", + Did: subjectDID, + Feature: feature, + CreatedAt: "2026-01-01T00:00:00Z", + }, aturi)) } func TestAllowVODUpload_InviteMode(t *testing.T) { diff --git a/pkg/spxrpc/labels_test.go b/pkg/spxrpc/labels_test.go index 3660923e..c01f2dc4 100644 --- a/pkg/spxrpc/labels_test.go +++ b/pkg/spxrpc/labels_test.go @@ -9,12 +9,14 @@ import ( "time" comatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/labstack/echo/v4" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/streamplace" "stream.place/streamplace/pkg/vod" ) @@ -113,13 +115,19 @@ func writeBlob(t *testing.T, store blob.Store, cid string) { func setupBlobTest(t *testing.T) (*Server, model.Model) { t.Helper() m := newTestModel(t) - require.NoError(t, m.UpsertMediaTrack(context.Background(), &model.MediaTrack{ - URI: "at://" + testOwner + "/place.stream.media.track/1", - RepoDID: testOwner, - Blob: testContentCID, - TrackID: "1", - MediaType: "video", - })) + aturi, err := syntax.ParseATURI("at://" + testOwner + "/place.stream.media.track/1") + require.NoError(t, err) + require.NoError(t, m.UpsertMediaTrack(context.Background(), &streamplace.MediaTrack{ + LexiconTypeID: "place.stream.media.track", + Track: &streamplace.MediaTrack_Track{ + MediaDefs_MuxlTrack: &streamplace.MediaDefs_MuxlTrack{ + LexiconTypeID: "place.stream.media.defs#muxlTrack", + Blob: testContentCID, + TrackId: "1", + MediaType: "video", + }, + }, + }, aturi)) store, err := blob.NewFileStore(t.TempDir()) require.NoError(t, err) writeBlob(t, store, testContentCID) diff --git a/pkg/spxrpc/place_stream_playback_getvideo.go b/pkg/spxrpc/place_stream_playback_getvideo.go index 5522e77f..8bdafae4 100644 --- a/pkg/spxrpc/place_stream_playback_getvideo.go +++ b/pkg/spxrpc/place_stream_playback_getvideo.go @@ -13,7 +13,6 @@ import ( "strings" "github.com/bluesky-social/indigo/atproto/syntax" - lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/labstack/echo/v4" "stream.place/streamplace/pkg/blob" @@ -432,27 +431,17 @@ func (s *Server) resolveVideoBlob(ctx context.Context, uri string) (*resolvedVid } // loadVideoRecord pulls the place.stream.video record at uri out of -// the local index and CBOR-decodes it. Returns a typed error suitable -// for surfacing back to the client (404, 422, 500). +// the local index. Returns a typed error suitable for surfacing back +// to the client (404, 500). func (s *Server) loadVideoRecord(ctx context.Context, uri string) (*streamplace.Video, error) { - video, err := s.model.GetVideoByURI(ctx, uri) + videoRec, err := s.model.GetVideoByURI(ctx, uri) if err != nil { log.Error(ctx, "playback: GetVideoByURI failed", "uri", uri, "error", err) return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) } - if video == nil { + if videoRec == nil { return nil, echo.NewHTTPError(http.StatusNotFound, "VideoNotFound") } - rec, err := lexutil.CborDecodeValue(video.Record) - if err != nil { - log.Error(ctx, "playback: decode video record failed", "uri", uri, "error", err) - return nil, echo.NewHTTPError(http.StatusInternalServerError, err.Error()) - } - videoRec, ok := rec.(*streamplace.Video) - if !ok { - return nil, echo.NewHTTPError(http.StatusInternalServerError, - fmt.Sprintf("video record at %s decoded as %T, expected *streamplace.Video", uri, rec)) - } return videoRec, nil } @@ -472,10 +461,14 @@ func (s *Server) firstTrackBlob(ctx context.Context, src *streamplace.MediaDefs_ log.Error(ctx, "playback: GetMediaTrackByURI failed", "uri", firstRef.Uri, "error", err) return "", echo.NewHTTPError(http.StatusInternalServerError, err.Error()) } - if track == nil || track.Blob == "" { + if track == nil || track.Track == nil || track.Track.MediaDefs_MuxlTrack == nil { + return "", echo.NewHTTPError(http.StatusNotFound, "TrackNotFound") + } + blob := track.Track.MediaDefs_MuxlTrack.Blob + if blob == "" { return "", echo.NewHTTPError(http.StatusNotFound, "TrackNotFound") } - return track.Blob, nil + return blob, nil } // fetchMetafile reads blobs/.json from the playback store and diff --git a/pkg/spxrpc/place_stream_playback_getvideo_test.go b/pkg/spxrpc/place_stream_playback_getvideo_test.go index 56107558..be7dfca6 100644 --- a/pkg/spxrpc/place_stream_playback_getvideo_test.go +++ b/pkg/spxrpc/place_stream_playback_getvideo_test.go @@ -1,14 +1,12 @@ package spxrpc import ( - "bytes" "context" "net/http" "net/http/httptest" "net/url" "strings" "testing" - "time" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" @@ -240,16 +238,6 @@ func TestSessionIDOrNew(t *testing.T) { }) } -// videoRecordCBOR marshals a *streamplace.Video to CBOR for stuffing -// into model.Video.Record so resolveVideoBlob's decode path can -// round-trip it back. -func videoRecordCBOR(t *testing.T, v *streamplace.Video) []byte { - t.Helper() - var buf bytes.Buffer - require.NoError(t, v.MarshalCBOR(&buf)) - return buf.Bytes() -} - // TestResolveVideoBlob_SourceClip walks the playback-resolve path for // a clip record: the clip's `Video` field points at a parent video // whose sourceTracks lead to a real MediaTrack + blob CID. resolve @@ -267,14 +255,23 @@ func TestResolveVideoBlob_SourceClip(t *testing.T) { clipURI = "at://did:plc:owner/place.stream.video/clipof" parentCID = "bafyparentblob" ) + parseURI := func(s string) syntax.ATURI { + u, err := syntax.ParseATURI(s) + require.NoError(t, err) + return u + } - require.NoError(t, m.UpsertMediaTrack(ctx, &model.MediaTrack{ - URI: trackURI, - RepoDID: owner, - Blob: parentCID, - TrackID: "1", - MediaType: "video", - })) + require.NoError(t, m.UpsertMediaTrack(ctx, &streamplace.MediaTrack{ + LexiconTypeID: "place.stream.media.track", + Track: &streamplace.MediaTrack_Track{ + MediaDefs_MuxlTrack: &streamplace.MediaDefs_MuxlTrack{ + LexiconTypeID: "place.stream.media.defs#muxlTrack", + Blob: parentCID, + TrackId: "1", + MediaType: "video", + }, + }, + }, parseURI(trackURI))) parent := &streamplace.Video{ LexiconTypeID: "place.stream.video", @@ -288,15 +285,7 @@ func TestResolveVideoBlob_SourceClip(t *testing.T) { }, }, } - require.NoError(t, m.UpsertVideo(ctx, &model.Video{ - URI: parentURI, - CID: "bafyparentrec", - RepoDID: owner, - RKey: "parent", - Title: parent.Title, - Record: videoRecordCBOR(t, parent), - IndexedAt: time.Now().UTC(), - })) + require.NoError(t, m.UpsertVideo(ctx, parent, parseURI(parentURI))) clip := &streamplace.Video{ LexiconTypeID: "place.stream.video", @@ -310,15 +299,7 @@ func TestResolveVideoBlob_SourceClip(t *testing.T) { }, }, } - require.NoError(t, m.UpsertVideo(ctx, &model.Video{ - URI: clipURI, - CID: "bafycliprec", - RepoDID: owner, - RKey: "clipof", - Title: clip.Title, - Record: videoRecordCBOR(t, clip), - IndexedAt: time.Now().UTC(), - })) + require.NoError(t, m.UpsertVideo(ctx, clip, parseURI(clipURI))) t.Run("non-clip parent resolves with no bounds", func(t *testing.T) { got, err := s.resolveVideoBlob(ctx, parentURI) @@ -354,15 +335,7 @@ func TestResolveVideoBlob_SourceClip(t *testing.T) { }, }, } - require.NoError(t, m.UpsertVideo(ctx, &model.Video{ - URI: nestedURI, - CID: "bafynested", - RepoDID: owner, - RKey: "nested", - Title: nested.Title, - Record: videoRecordCBOR(t, nested), - IndexedAt: time.Now().UTC(), - })) + require.NoError(t, m.UpsertVideo(ctx, nested, parseURI(nestedURI))) _, err := s.resolveVideoBlob(ctx, nestedURI) require.Error(t, err) he, ok := err.(*echo.HTTPError) -- 2.51.2