diff --git a/js/app/src/router.tsx b/js/app/src/router.tsx index 290a38df0..055b7a4f3 100644 --- a/js/app/src/router.tsx +++ b/js/app/src/router.tsx @@ -543,16 +543,20 @@ export function StreamplaceDrawer() { const [isLiveDashboard, setIsLiveDashboard] = useState(false); useEffect(() => { if (!isLiveDashboard && userIsLive) { - toast.show("You are live!", "Do you want to go to your Live Dashboard?", { - actionLabel: "Go", - onAction: () => { - navigation.navigate("LiveDashboard"); - setLivePopup(false); + toast.show( + "You are streaming!", + "Do you want to go to your Live Dashboard?", + { + actionLabel: "Go", + onAction: () => { + navigation.navigate("LiveDashboard"); + setLivePopup(false); + }, + onClose: () => setLivePopup(false), + variant: "error", + duration: 8, }, - onClose: () => setLivePopup(false), - variant: "error", - duration: 8, - }); + ); } }, [userIsLive]); const externalItems = useExternalItems(); diff --git a/pkg/localdb/localdb.go b/pkg/localdb/localdb.go index 0b186bc08..64ee5b151 100644 --- a/pkg/localdb/localdb.go +++ b/pkg/localdb/localdb.go @@ -17,7 +17,7 @@ type LocalDB interface { CreateSegment(segment *Segment) error MostRecentSegments() ([]Segment, error) LatestSegmentForUser(user string) (*Segment, error) - LatestSegmentsForUser(user string, limit int, before *time.Time, after *time.Time) ([]Segment, error) + LatestSegmentsForUser(user string, limit int, includeUnpublished bool, before *time.Time, after *time.Time) ([]Segment, error) FilterLiveRepoDIDs(repoDIDs []string) ([]string, error) CreateThumbnail(thumb *Thumbnail) error LatestThumbnailForUser(user string) (*Thumbnail, error) diff --git a/pkg/localdb/segment.go b/pkg/localdb/segment.go index 3cfc810cd..d8d8b667b 100644 --- a/pkg/localdb/segment.go +++ b/pkg/localdb/segment.go @@ -138,8 +138,8 @@ func (c ContentWarningsSlice) Value() (driver.Value, error) { type Segment struct { ID string `json:"id" gorm:"primaryKey"` SigningKeyDID string `json:"signingKeyDID" gorm:"column:signing_key_did"` - StartTime time.Time `json:"startTime" gorm:"index:latest_segments,priority:2;index:start_time"` - RepoDID string `json:"repoDID" gorm:"index:latest_segments,priority:1;column:repo_did"` + StartTime time.Time `json:"startTime" gorm:"index:latest_segments,priority:2;index:start_time;index:latest_segments_published,priority:2"` + RepoDID string `json:"repoDID" gorm:"index:latest_segments,priority:1;column:repo_did;index:latest_segments_published,priority:1"` Title string `json:"title"` Size int `json:"size" gorm:"column:size"` MediaData *SegmentMediaData `json:"mediaData,omitempty"` @@ -147,6 +147,7 @@ type Segment struct { ContentRights *ContentRights `json:"contentRights,omitempty"` DistributionPolicy *DistributionPolicy `json:"distributionPolicy,omitempty"` DeleteAfter *time.Time `json:"deleteAfter,omitempty" gorm:"column:delete_after;index:delete_after"` + Published bool `json:"published" gorm:"column:published;index:latest_segments_published,priority:3"` } func (s *Segment) ToStreamplaceSegment() (*streamplace.Segment, error) { @@ -238,7 +239,7 @@ func (m *LocalDatabase) MostRecentSegments() ([]Segment, error) { err := m.DB.Table("segments"). Select("segments.*"). - Where("start_time > ?", thirtySecondsAgo.UTC()). + Where("start_time > ? AND published = ?", thirtySecondsAgo.UTC(), true). Order("start_time DESC"). Find(&segments).Error if err != nil { @@ -270,7 +271,7 @@ func (m *LocalDatabase) MostRecentSegments() ([]Segment, error) { func (m *LocalDatabase) LatestSegmentForUser(user string) (*Segment, error) { var seg Segment - err := m.DB.Model(Segment{}).Where("repo_did = ?", user).Order("start_time DESC").First(&seg).Error + err := m.DB.Model(Segment{}).Where("repo_did = ? AND published = ?", user, true).Order("start_time DESC").First(&seg).Error if err != nil { return nil, err } @@ -288,7 +289,7 @@ func (m *LocalDatabase) FilterLiveRepoDIDs(repoDIDs []string) ([]string, error) err := m.DB.Table("segments"). Select("DISTINCT repo_did"). - Where("repo_did IN ? AND start_time > ?", repoDIDs, thirtySecondsAgo.UTC()). + Where("repo_did IN ? AND start_time > ? AND published = ?", repoDIDs, thirtySecondsAgo.UTC(), true). Pluck("repo_did", &liveDIDs).Error if err != nil { @@ -298,7 +299,7 @@ func (m *LocalDatabase) FilterLiveRepoDIDs(repoDIDs []string) ([]string, error) return liveDIDs, nil } -func (m *LocalDatabase) LatestSegmentsForUser(user string, limit int, before *time.Time, after *time.Time) ([]Segment, error) { +func (m *LocalDatabase) LatestSegmentsForUser(user string, limit int, includeUnpublished bool, before *time.Time, after *time.Time) ([]Segment, error) { var segs []Segment if before == nil { later := time.Now().Add(1000 * time.Hour) @@ -308,7 +309,11 @@ func (m *LocalDatabase) LatestSegmentsForUser(user string, limit int, before *ti earlier := time.Time{} after = &earlier } - err := m.DB.Model(Segment{}).Where("repo_did = ? AND start_time < ? AND start_time > ?", user, before.UTC(), after.UTC()).Order("start_time DESC").Limit(limit).Find(&segs).Error + query := m.DB.Model(Segment{}).Where("repo_did = ? AND start_time < ? AND start_time > ?", user, before.UTC(), after.UTC()) + if !includeUnpublished { + query = query.Where("published = ?", true) + } + err := query.Order("start_time DESC").Limit(limit).Find(&segs).Error if err != nil { return nil, err } diff --git a/pkg/media/clip_user.go b/pkg/media/clip_user.go index 51f34f3ea..fe365c448 100644 --- a/pkg/media/clip_user.go +++ b/pkg/media/clip_user.go @@ -14,7 +14,7 @@ import ( ) func ClipUser(ctx context.Context, localDB localdb.LocalDB, cli *config.CLI, user string, writer io.Writer, before *time.Time, after *time.Time) error { - segments, err := localDB.LatestSegmentsForUser(user, -1, before, after) + segments, err := localDB.LatestSegmentsForUser(user, -1, false, before, after) if err != nil { return fmt.Errorf("unable to get segments: %w", err) } diff --git a/pkg/media/validate.go b/pkg/media/validate.go index f28bae313..871ab4ed2 100644 --- a/pkg/media/validate.go +++ b/pkg/media/validate.go @@ -129,6 +129,7 @@ func (mm *MediaManager) ValidateMP4(ctx context.Context, input io.Reader, local ContentRights: meta.ContentRights, DistributionPolicy: meta.DistributionPolicy, DeleteAfter: deleteAfter, + Published: meta.Published, } mm.newSegmentSubsMutex.RLock() defer mm.newSegmentSubsMutex.RUnlock() diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index a84f0b0af..9d6e01693 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -82,7 +82,10 @@ func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before st beforeTime = &parsedTime } - segments, err := s.localDB.LatestSegmentsForUser(userDID, limit, beforeTime, nil) + sess, _ := oatproxy.GetOAuthSession(ctx) + includeUnpublished := sess != nil && sess.DID == userDID + + segments, err := s.localDB.LatestSegmentsForUser(userDID, limit, includeUnpublished, beforeTime, nil) if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to fetch segments") }