diff --git a/pkg/api/playback.go b/pkg/api/playback.go index 4b813f74a..b31dd8eec 100644 --- a/pkg/api/playback.go +++ b/pkg/api/playback.go @@ -11,12 +11,10 @@ import ( "github.com/julienschmidt/httprouter" "github.com/pion/webrtc/v4" - "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/errors" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/spmetrics" - "stream.place/streamplace/pkg/streamplace" ) func (a *StreamplaceAPI) NormalizeUser(ctx context.Context, user string) (string, error) { @@ -261,6 +259,11 @@ func (a *StreamplaceAPI) HandleHLSPlayback(ctx context.Context) httprouter.Handl }) } +// thumbnailMaxAge is how stale a thumbnail may be before we treat the user as +// offline and stop serving it. It must comfortably exceed thumbnailInterval (the +// rate at which live thumbnails are refreshed) to avoid flickering mid-stream. +const thumbnailMaxAge = 2 * time.Minute + func (a *StreamplaceAPI) HandleThumbnailPlayback(ctx context.Context) httprouter.Handle { return func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { ctx = log.WithLogValues(r.Context(), "func", "HandleThumbnailPlayback") @@ -274,44 +277,16 @@ func (a *StreamplaceAPI) HandleThumbnailPlayback(ctx context.Context) httprouter errors.WriteHTTPNotFound(w, "user not found", err) return } - if !a.CLI.WideOpen { - ls, err := a.Model.GetLatestLivestreamForRepo(user) - if err != nil { - errors.WriteHTTPInternalServerError(w, "could not get livestream", err) - return - } - if ls == nil { - errors.WriteHTTPNotFound(w, "livestream not found", err) - return - } - lsrv, err := ls.ToLivestreamView() - if err != nil { - errors.WriteHTTPInternalServerError(w, "could not marshal livestream", err) - return - } - lsr, ok := lsrv.Record.Val.(*streamplace.Livestream) - if !ok { - errors.WriteHTTPInternalServerError(w, "livestream is not a streamplace livestream", nil) - return - } - if lsr.EndedAt != nil { - errors.WriteHTTPNotFound(w, "livestream has ended", nil) - return - } - } - thumb, err := a.LocalDB.LatestThumbnailForUser(user) - if err != nil { - errors.WriteHTTPInternalServerError(w, "could not query thumbnail", err) - return - } - if thumb == nil { - errors.WriteHTTPNotFound(w, "thumbnail not found", err) + fpath := a.CLI.ThumbnailFilePath(user) + mt, ok := a.CLI.ThumbnailModTime(user) + if !ok { + errors.WriteHTTPNotFound(w, "thumbnail not found", nil) return } - aqt := aqtime.FromTime(thumb.Segment.StartTime) - fpath, err := a.CLI.SegmentFilePath(user, fmt.Sprintf("%s.%s", aqt.String(), thumb.Format)) - if err != nil { - errors.WriteHTTPInternalServerError(w, "could not get segment file path", err) + // A thumbnail that hasn't been refreshed recently means the user is no + // longer live, so don't serve a stale preview. WideOpen (dev) skips this. + if !a.CLI.WideOpen && time.Since(mt) > thumbnailMaxAge { + errors.WriteHTTPNotFound(w, "no recent thumbnail", nil) return } log.Debug(ctx, "serving thumbnail", "fpath", fpath) diff --git a/pkg/config/config.go b/pkg/config/config.go index 412b60032..7911510a3 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -37,6 +37,7 @@ import ( const SPDataDir = "$SP_DATA_DIR" const SegmentsDir = "segments" +const ThumbnailsDir = "thumbnails" type BuildFlags struct { Version string @@ -1269,6 +1270,46 @@ func (cli *CLI) SegmentFileCreate(user string, aqt aqtime.AQTime, ext string) (* return cli.DataFileCreate([]string{SegmentsDir, user, yr, mon, day, hr, min, fname}, false) } +// ThumbnailFilePath returns the path to a user's current thumbnail. There is a +// single, continually-overwritten thumbnail per user. +func (cli *CLI) ThumbnailFilePath(user string) string { + return cli.DataFilePath([]string{ThumbnailsDir, fmt.Sprintf("%s.jpg", user)}) +} + +// ThumbnailModTime returns the modification time of a user's thumbnail and +// whether it exists. The mod time doubles as a "last seen live" signal. +func (cli *CLI) ThumbnailModTime(user string) (time.Time, bool) { + fi, err := os.Stat(cli.ThumbnailFilePath(user)) + if err != nil { + return time.Time{}, false + } + return fi.ModTime(), true +} + +// ThumbnailWrite atomically (re)writes a user's thumbnail. The image is written +// to a temp file via the supplied function and renamed into place, so readers +// (and PDS uploads) never observe a half-written thumbnail. +func (cli *CLI) ThumbnailWrite(user string, write func(io.Writer) error) error { + final := cli.ThumbnailFilePath(user) + dir := filepath.Dir(final) + if err := os.MkdirAll(dir, os.ModePerm); err != nil { + return fmt.Errorf("error creating thumbnail dir %s: %w", dir, err) + } + tmp, err := os.CreateTemp(dir, "thumb-*.jpg") + if err != nil { + return err + } + defer os.Remove(tmp.Name()) // no-op once the rename below succeeds + if err := write(tmp); err != nil { + tmp.Close() + return err + } + if err := tmp.Close(); err != nil { + return err + } + return os.Rename(tmp.Name(), final) +} + // read a file from our data dir func (cli *CLI) DataFileRead(fpath []string, w io.Writer) error { ddpath := cli.DataFilePath(fpath) diff --git a/pkg/config/thumbnail_test.go b/pkg/config/thumbnail_test.go new file mode 100644 index 000000000..37473f3d4 --- /dev/null +++ b/pkg/config/thumbnail_test.go @@ -0,0 +1,67 @@ +package config + +import ( + "fmt" + "io" + "os" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestThumbnailWriteAndRead(t *testing.T) { + cli := &CLI{DataDir: t.TempDir()} + const user = "did:plc:abc123" + + // No thumbnail exists yet. + _, ok := cli.ThumbnailModTime(user) + require.False(t, ok) + + require.NoError(t, cli.ThumbnailWrite(user, func(w io.Writer) error { + _, err := io.WriteString(w, "first") + return err + })) + + // Path lives under the thumbnails dir, colon-sanitized, with a .jpg suffix. + fpath := cli.ThumbnailFilePath(user) + require.Equal(t, filepath.Join(cli.DataDir, ThumbnailsDir, "did-plc-abc123.jpg"), fpath) + + data, err := os.ReadFile(fpath) + require.NoError(t, err) + require.Equal(t, "first", string(data)) + + mt, ok := cli.ThumbnailModTime(user) + require.True(t, ok) + require.WithinDuration(t, time.Now(), mt, 5*time.Second) + + // A second write overwrites in place... + require.NoError(t, cli.ThumbnailWrite(user, func(w io.Writer) error { + _, err := io.WriteString(w, "second") + return err + })) + data, err = os.ReadFile(fpath) + require.NoError(t, err) + require.Equal(t, "second", string(data)) + + // ...and leaves no temp files behind. + requireOnlyThumbnail(t, cli, "did-plc-abc123.jpg") + + // A failed write keeps the previous thumbnail intact and cleans up its temp. + require.Error(t, cli.ThumbnailWrite(user, func(w io.Writer) error { + return fmt.Errorf("boom") + })) + data, err = os.ReadFile(fpath) + require.NoError(t, err) + require.Equal(t, "second", string(data)) + requireOnlyThumbnail(t, cli, "did-plc-abc123.jpg") +} + +func requireOnlyThumbnail(t *testing.T, cli *CLI, name string) { + t.Helper() + entries, err := os.ReadDir(filepath.Join(cli.DataDir, ThumbnailsDir)) + require.NoError(t, err) + require.Len(t, entries, 1) + require.Equal(t, name, entries[0].Name()) +} diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 9f797339b..5cdb4b1a7 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "fmt" + "io" "net/url" "time" @@ -310,42 +311,27 @@ func shouldNotify(lsv *streamplace.Livestream_LivestreamView) bool { return *settings.PushNotification } +// thumbnailInterval is how often we refresh a user's thumbnail while they're +// live. A missing or older thumbnail (e.g. the user just went live) is +// regenerated immediately on the next segment. +const thumbnailInterval = 30 * time.Second + func (ss *StreamSession) Thumbnail(ctx context.Context, repoDID string, not *media.NewSegmentNotification) error { - lock := thumbnail.GetThumbnailLock(not.Segment.RepoDID) - locked := lock.TryLock() - if !locked { + lock := thumbnail.GetThumbnailLock(repoDID) + if !lock.TryLock() { // we're already generating a thumbnail for this user, skip return nil } defer lock.Unlock() - oldThumb, err := ss.localDB.LatestThumbnailForUser(not.Segment.RepoDID) - if err != nil { - return err - } - if oldThumb != nil && not.Segment.StartTime.Sub(oldThumb.Segment.StartTime) < time.Minute { - // we have a thumbnail <60sec old, skip generating a new one + + if mt, ok := ss.cli.ThumbnailModTime(repoDID); ok && time.Since(mt) < thumbnailInterval { + // current thumbnail is still fresh, keep it return nil } - r := bytes.NewReader(not.Data) - aqt := aqtime.FromTime(not.Segment.StartTime) - fd, err := ss.cli.SegmentFileCreate(not.Segment.RepoDID, aqt, "jpeg") - if err != nil { - return err - } - defer fd.Close() - err = media.Thumbnail(ctx, r, fd, "jpeg") - if err != nil { - return err - } - thumb := &localdb.Thumbnail{ - Format: "jpeg", - SegmentID: not.Segment.ID, - } - err = ss.localDB.CreateThumbnail(thumb) - if err != nil { - return err - } - return nil + + return ss.cli.ThumbnailWrite(repoDID, func(w io.Writer) error { + return media.Thumbnail(ctx, bytes.NewReader(not.Data), w, "jpeg") + }) } // UpdateStatus signals the background worker to update status (non-blocking) diff --git a/pkg/localdb/localdb.go b/pkg/localdb/localdb.go index 9a112e270..1e180b3ca 100644 --- a/pkg/localdb/localdb.go +++ b/pkg/localdb/localdb.go @@ -19,9 +19,6 @@ type LocalDB interface { LatestSegmentForUser(user string) (*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) - ThumbnailCleaner(ctx context.Context) error GetSegment(id string) (*Segment, error) GetExpiredSegments(ctx context.Context) ([]Segment, error) DeleteSegment(ctx context.Context, id string) error @@ -74,7 +71,6 @@ func MakeDB(dbURL string) (LocalDB, error) { sqlDB.SetMaxOpenConns(1) for _, model := range []any{ Segment{}, - Thumbnail{}, ViewLogSalt{}, } { err = db.AutoMigrate(model) @@ -82,5 +78,13 @@ func MakeDB(dbURL string) (LocalDB, error) { return nil, err } } + // Thumbnails used to live in this table but are now served from the + // filesystem (see config.ThumbnailFilePath). Drop the legacy table so it + // stops bloating the database. + if db.Migrator().HasTable("thumbnails") { + if err := db.Migrator().DropTable("thumbnails"); err != nil { + return nil, fmt.Errorf("error dropping legacy thumbnails table: %w", err) + } + } return &LocalDatabase{DB: db}, nil } diff --git a/pkg/localdb/thumbnail.go b/pkg/localdb/thumbnail.go deleted file mode 100644 index 56ed3cd9e..000000000 --- a/pkg/localdb/thumbnail.go +++ /dev/null @@ -1,124 +0,0 @@ -package localdb - -import ( - "context" - "fmt" - - "github.com/google/uuid" - "stream.place/streamplace/pkg/log" -) - -type Thumbnail struct { - ID string `json:"id" gorm:"primaryKey"` - Format string `json:"format"` - SegmentID string `json:"segmentId" gorm:"index"` - Segment Segment `json:"segment,omitempty" gorm:"foreignKey:SegmentID;references:id"` -} - -func (m *LocalDatabase) CreateThumbnail(thumb *Thumbnail) error { - uu, err := uuid.NewV7() - if err != nil { - return err - } - if thumb.SegmentID == "" { - return fmt.Errorf("segmentID is required") - } - thumb.ID = uu.String() - err = m.DB.Model(Thumbnail{}).Create(thumb).Error - if err != nil { - return err - } - return nil -} - -// return the most recent thumbnail for a user -func (m *LocalDatabase) LatestThumbnailForUser(user string) (*Thumbnail, error) { - var thumbnail Thumbnail - - res := m.DB.Table("thumbnails AS t"). - Select("t.*"). - Joins("JOIN segments AS s ON t.segment_id = s.id"). - Where("s.repo_did = ?", user). - Order("s.start_time DESC"). - Limit(1). - Scan(&thumbnail) - - if res.RowsAffected == 0 { - return nil, nil - } - if res.Error != nil { - return nil, res.Error - } - - var seg Segment - err := m.DB.First(&seg, "id = ?", thumbnail.SegmentID).Error - if err != nil { - return nil, fmt.Errorf("could not find segment for thumbnail SegmentID=%s", thumbnail.SegmentID) - } - - thumbnail.Segment = seg - - return &thumbnail, nil -} - -// ThumbnailCleaner keeps only the most recent thumbnail for each user and -// deletes the rest. Thumbnails are created roughly once a minute per active -// user but are never removed when their segment is cleaned up, so the table -// grows without bound and bloats the database. Only the latest thumbnail per -// user is ever read (see LatestThumbnailForUser), so the older ones are dead -// weight. -func (m *LocalDatabase) ThumbnailCleaner(ctx context.Context) error { - var cleaned int64 - - // Drop any thumbnail whose segment has already been deleted. These are - // orphans we can never serve and can't even associate with a user anymore. - orphans := m.DB. - Where("segment_id NOT IN (?)", m.DB.Model(&Segment{}).Select("id")). - Delete(&Thumbnail{}) - if orphans.Error != nil { - log.Error(ctx, "Failed to clean orphaned thumbnails", "error", orphans.Error) - return orphans.Error - } - cleaned += orphans.RowsAffected - - // Find all unique repo_did values. - var repoDIDs []string - if err := m.DB.Model(&Segment{}).Distinct("repo_did").Pluck("repo_did", &repoDIDs).Error; err != nil { - log.Error(ctx, "Failed to get unique repo_dids for thumbnail cleaning", "error", err) - return err - } - - // For each user, keep the thumbnail on their most recent segment (matching - // LatestThumbnailForUser) and delete every other thumbnail of theirs. - for _, repoDID := range repoDIDs { - var keepIDs []string - if err := m.DB.Table("thumbnails AS t"). - Joins("JOIN segments AS s ON t.segment_id = s.id"). - Where("s.repo_did = ?", repoDID). - Order("s.start_time DESC"). - Limit(1). - Pluck("t.id", &keepIDs).Error; err != nil { - log.Error(ctx, "Failed to get thumbnail to keep", "repo_did", repoDID, "error", err) - return err - } - if len(keepIDs) == 0 { - continue - } - - result := m.DB. - Where("segment_id IN (?) AND id NOT IN ?", - m.DB.Model(&Segment{}).Select("id").Where("repo_did = ?", repoDID), - keepIDs). - Delete(&Thumbnail{}) - if result.Error != nil { - log.Error(ctx, "Failed to clean old thumbnails", "repo_did", repoDID, "error", result.Error) - return result.Error - } - cleaned += result.RowsAffected - } - - if cleaned > 0 { - log.Log(ctx, "Cleaned old thumbnails", "count", cleaned) - } - return nil -} diff --git a/pkg/localdb/thumbnail_test.go b/pkg/localdb/thumbnail_test.go deleted file mode 100644 index 38df99069..000000000 --- a/pkg/localdb/thumbnail_test.go +++ /dev/null @@ -1,84 +0,0 @@ -package localdb - -import ( - "context" - "testing" - "time" - - "github.com/stretchr/testify/require" - "stream.place/streamplace/pkg/config" -) - -func TestThumbnailCleaner(t *testing.T) { - config.DisableSQLLogging() - defer config.EnableSQLLogging() - - db, err := MakeDB(":memory:") - require.NoError(t, err) - ldb := db.(*LocalDatabase) - - const ( - userA = "did:plc:aaa" - userB = "did:plc:bbb" - userC = "did:plc:ccc" - ) - base := time.Now().UTC() - - mkSeg := func(id, user string, startOffset time.Duration) { - require.NoError(t, ldb.CreateSegment(&Segment{ - ID: id, - RepoDID: user, - StartTime: base.Add(startOffset), - Published: true, - })) - } - mkThumb := func(segID string) { - require.NoError(t, ldb.CreateThumbnail(&Thumbnail{Format: "jpeg", SegmentID: segID})) - } - countThumbs := func() int64 { - var n int64 - require.NoError(t, ldb.DB.Model(&Thumbnail{}).Count(&n).Error) - return n - } - - // userA: three segments, each with a thumbnail; a3 is the most recent. - mkSeg("a1", userA, -3*time.Minute) - mkThumb("a1") - mkSeg("a2", userA, -2*time.Minute) - mkThumb("a2") - mkSeg("a3", userA, -1*time.Minute) - mkThumb("a3") - - // userB: two segments with thumbnails; b2 is the most recent. - mkSeg("b1", userB, -5*time.Minute) - mkThumb("b1") - mkSeg("b2", userB, -4*time.Minute) - mkThumb("b2") - - // userC: has a segment but no thumbnail (must not break anything). - mkSeg("c1", userC, -1*time.Minute) - - // An orphaned thumbnail whose segment was already cleaned up. - mkThumb("ghost") - - require.EqualValues(t, 6, countThumbs()) - - require.NoError(t, ldb.ThumbnailCleaner(context.Background())) - - // One thumbnail per user with thumbnails should remain; orphan is gone. - require.EqualValues(t, 2, countThumbs()) - - ta, err := ldb.LatestThumbnailForUser(userA) - require.NoError(t, err) - require.NotNil(t, ta) - require.Equal(t, "a3", ta.SegmentID) - - tb, err := ldb.LatestThumbnailForUser(userB) - require.NoError(t, err) - require.NotNil(t, tb) - require.Equal(t, "b2", tb.SegmentID) - - // Running again is a no-op. - require.NoError(t, ldb.ThumbnailCleaner(context.Background())) - require.EqualValues(t, 2, countThumbs()) -} diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index 3802a75e9..8c703052a 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -20,7 +20,6 @@ import ( "github.com/gorilla/websocket" "github.com/labstack/echo/v4" "github.com/streamplace/oatproxy/pkg/oatproxy" - "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/spid" "stream.place/streamplace/pkg/spmetrics" @@ -421,36 +420,18 @@ func (s *Server) handlePlaceStreamLiveStartLivestream(ctx context.Context, body livestream.LastSeenAt = &now if livestream.Thumb == nil { - // Step 1: get latest thumbnail from localDB and upload to user's PDS + // Upload the user's current thumbnail to their PDS as the livestream image. var thumb *lexutil.LexBlob - dbThumb, err := s.localDB.LatestThumbnailForUser(session.DID) + thumbData, err := os.ReadFile(s.cli.ThumbnailFilePath(session.DID)) if err != nil { - log.Error(ctx, "failed to get latest thumbnail", "err", err) - } - if dbThumb != nil { - aqt := aqtime.FromTime(dbThumb.Segment.StartTime) - fpath, err := s.cli.SegmentFilePath(session.DID, fmt.Sprintf("%s.%s", aqt.String(), dbThumb.Format)) + log.Error(ctx, "failed to read thumbnail file", "err", err) + } else { + var uploadOut comatproto.RepoUploadBlob_Output + err = client.Do(ctx, xrpc.Procedure, "image/jpeg", "com.atproto.repo.uploadBlob", nil, bytes.NewReader(thumbData), &uploadOut) if err != nil { - log.Error(ctx, "failed to get thumbnail file path", "err", err) + log.Error(ctx, "failed to upload thumbnail to PDS", "err", err) } else { - thumbData, err := os.ReadFile(fpath) - if err != nil { - log.Error(ctx, "failed to read thumbnail file", "err", err) - } else { - mimeType := "image/jpeg" - if dbThumb.Format == "png" { - mimeType = "image/png" - } - - // Step 2: upload to user's PDS - var uploadOut comatproto.RepoUploadBlob_Output - err = client.Do(ctx, xrpc.Procedure, mimeType, "com.atproto.repo.uploadBlob", nil, bytes.NewReader(thumbData), &uploadOut) - if err != nil { - log.Error(ctx, "failed to upload thumbnail to PDS", "err", err) - } else { - thumb = uploadOut.Blob - } - } + thumb = uploadOut.Blob } } livestream.Thumb = thumb diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go index 1bcd2ec3f..15535071e 100644 --- a/pkg/storage/storage.go +++ b/pkg/storage/storage.go @@ -25,12 +25,6 @@ func StartSegmentCleaner(ctx context.Context, localDB localdb.LocalDB, cli *conf case <-ctx.Done(): return nil case <-time.After(60 * time.Second): - // Keep only the latest thumbnail per user; the rest accumulate - // forever and bloat the database. Log and continue on error so a - // thumbnail hiccup can't take down the segment cleaner. - if err := localDB.ThumbnailCleaner(ctx); err != nil { - log.Error(ctx, "Failed to clean thumbnails", "error", err) - } expiredSegments, err := localDB.GetExpiredSegments(ctx) if err != nil { return err