From 1a1a5f810a8314c8e9d4862efdb3c8a3b4387679 Mon Sep 17 00:00:00 2001 From: zijiren233 Date: Wed, 18 Oct 2023 19:33:35 +0800 Subject: [PATCH] Feat: lazy load movies --- internal/bootstrap/room.go | 14 +---- internal/bootstrap/rtmp.go | 37 ++++++++++++- internal/conf/rtmp.go | 2 - internal/db/movie.go | 4 ++ internal/op/movie.go | 16 +++++- internal/op/room.go | 111 ++++++++++++++++++++++++------------- internal/op/rooms.go | 1 + internal/rtmp/rtmp.go | 18 ------ server/handlers/movie.go | 14 ++--- server/handlers/public.go | 1 - server/handlers/room.go | 4 +- utils/utils.go | 4 ++ 12 files changed, 143 insertions(+), 83 deletions(-) diff --git a/internal/bootstrap/room.go b/internal/bootstrap/room.go index 7e67c464..410a4b78 100644 --- a/internal/bootstrap/room.go +++ b/internal/bootstrap/room.go @@ -14,23 +14,11 @@ func InitRoom(ctx context.Context) error { return err } for _, room := range r { - r, err := op.LoadRoom(room) + _, err := op.LoadRoom(room) if err != nil { log.Errorf("load room error: %v", err) return err } - m, err := r.GetAllMoviesByRoomID() - if err != nil { - log.Errorf("get all movies by room id error: %v", err) - return err - } - for _, movie := range m { - err = r.InitMovie(movie) - if err != nil { - log.Errorf("init movie error: %v", err) - return err - } - } } return nil } diff --git a/internal/bootstrap/rtmp.go b/internal/bootstrap/rtmp.go index 706be289..8ac756f4 100644 --- a/internal/bootstrap/rtmp.go +++ b/internal/bootstrap/rtmp.go @@ -2,14 +2,49 @@ package bootstrap import ( "context" + "fmt" + "strconv" + log "github.com/sirupsen/logrus" "github.com/synctv-org/synctv/internal/conf" + "github.com/synctv-org/synctv/internal/op" "github.com/synctv-org/synctv/internal/rtmp" rtmps "github.com/zijiren233/livelib/server" ) func InitRtmp(ctx context.Context) error { - s := rtmps.NewRtmpServer(rtmps.WithInitHlsPlayer(conf.Conf.Rtmp.HlsPlayer)) + s := rtmps.NewRtmpServer(rtmps.WithInitHlsPlayer(true)) rtmp.Init(s) + + s.SetParseChannelFunc(func(ReqAppName, ReqChannelName string, IsPublisher bool) (TrueAppName string, TrueChannel string, err error) { + if IsPublisher { + channelName, err := rtmp.AuthRtmpPublish(ReqChannelName) + if err != nil { + log.Errorf("rtmp: publish auth to %s error: %v", ReqAppName, err) + return "", "", err + } + log.Infof("rtmp: publisher login success: %s/%s", ReqAppName, channelName) + id, err := strconv.Atoi(ReqAppName) + if err != nil { + log.Errorf("rtmp: parse channel name to id error: %v", err) + return "", "", err + } + r, err := op.GetRoomByID(uint(id)) + if err != nil { + log.Errorf("rtmp: get room by id error: %v", err) + return "", "", err + } + err = r.LazyInit() + if err != nil { + log.Errorf("rtmp: lazy init room error: %v", err) + return "", "", err + } + return ReqAppName, channelName, nil + } else if !conf.Conf.Rtmp.RtmpPlayer { + log.Warnf("rtmp: dial to %s/%s error: %s", ReqAppName, ReqChannelName, "rtmp player is not enabled") + return "", "", fmt.Errorf("rtmp: dial to %s/%s error: %s", ReqAppName, ReqChannelName, "rtmp player is not enabled") + } + return ReqAppName, ReqChannelName, nil + }) return nil } diff --git a/internal/conf/rtmp.go b/internal/conf/rtmp.go index ef9d8737..220d6bbc 100644 --- a/internal/conf/rtmp.go +++ b/internal/conf/rtmp.go @@ -6,7 +6,6 @@ type RtmpConfig struct { CustomPublishHost string `yaml:"custom_publish_host" lc:"publish host (default use http header host)" env:"RTMP_CUSTOM_PUBLISH_HOST"` RtmpPlayer bool `yaml:"rtmp_player" lc:"enable rtmp player (default: false)" env:"RTMP_PLAYER"` - HlsPlayer bool `yaml:"hls_player" lc:"enable hls player (default: true)" env:"HLS_PLAYER"` } func DefaultRtmpConfig() RtmpConfig { @@ -15,6 +14,5 @@ func DefaultRtmpConfig() RtmpConfig { Port: 0, CustomPublishHost: "", RtmpPlayer: false, - HlsPlayer: true, } } diff --git a/internal/db/movie.go b/internal/db/movie.go index 13c94220..5a5ff80c 100644 --- a/internal/db/movie.go +++ b/internal/db/movie.go @@ -40,6 +40,10 @@ func UpdateMovie(movie *model.Movie, columns ...clause.Column) error { return db.Model(movie).Clauses(clause.Returning{Columns: columns}).Where("room_id = ? AND id = ?", movie.RoomID, movie.ID).Updates(movie).Error } +func SaveMovie(movie *model.Movie, columns ...clause.Column) error { + return db.Model(movie).Clauses(clause.Returning{Columns: columns}).Where("room_id = ? AND id = ?", movie.RoomID, movie.ID).Save(movie).Error +} + func SwapMoviePositions(roomID uint, movie1ID uint, movie2ID uint) (err error) { tx := db.Begin() defer func() { diff --git a/internal/op/movie.go b/internal/op/movie.go index c971487d..24f1b490 100644 --- a/internal/op/movie.go +++ b/internal/op/movie.go @@ -5,6 +5,7 @@ import ( "time" "github.com/bluele/gcache" + log "github.com/sirupsen/logrus" "github.com/synctv-org/synctv/internal/db" "github.com/synctv-org/synctv/internal/model" "github.com/zijiren233/gencontainer/dllist" @@ -17,7 +18,6 @@ var movieCache = gcache.New(2048). func GetAllMoviesByRoomID(roomID uint) (*dllist.Dllist[*model.Movie], error) { i, err := movieCache.Get(roomID) if err == nil { - return i.(*dllist.Dllist[*model.Movie]), nil } m, err := db.GetAllMoviesByRoomID(roomID) @@ -103,6 +103,20 @@ func UpdateMovie(movie *model.Movie) error { return nil } +func SaveMovie(movie *model.Movie) error { + log.Debug(movie) + err := db.SaveMovie(movie) + if err != nil { + return err + } + m, err := GetMovieByID(movie.RoomID, movie.ID) + if err != nil { + return err + } + *m = *movie + return nil +} + func DeleteMoviesByRoomID(roomID uint) error { movieCache.Remove(roomID) return db.DeleteMoviesByRoomID(roomID) diff --git a/internal/op/room.go b/internal/op/room.go index 5183786b..d250bf17 100644 --- a/internal/op/room.go +++ b/internal/op/room.go @@ -3,6 +3,7 @@ package op import ( "errors" "net/url" + "strconv" "sync/atomic" "time" @@ -34,26 +35,52 @@ type Room struct { hub *Hub } -func (r *Room) lazyInit() { +func (r *Room) LazyInit() (err error) { r.initOnce.Do(func() { - r.current = newCurrent() r.hub = newHub(r.ID) - a, err := rtmp.RtmpServer().NewApp(r.Name) + r.rtmpa, err = rtmp.RtmpServer().NewApp(strconv.Itoa(int(r.ID))) if err != nil { - log.Fatalf("failed to create rtmp app: %s", err.Error()) + log.Errorf("failed to create rtmp app: %s", err.Error()) + return + } + + var ms []*model.Movie + ms, err = r.GetAllMoviesByRoomID() + if err != nil { + log.Errorf("failed to get movies: %s", err.Error()) + return + } + for _, m := range ms { + if err = r.initMovie(m); err != nil { + log.Errorf("failed to init movie: %s", err.Error()) + return + } } - r.rtmpa = a }) + return } -func (r *Room) Hub() *Hub { - r.lazyInit() - return r.hub +func (r *Room) ClientNum() int64 { + if r.hub == nil { + return 0 + } + return r.hub.ClientNum() } -func (r *Room) App() *rtmps.App { - r.lazyInit() - return r.rtmpa +func (r *Room) Broadcast(data Message, conf ...BroadcastConf) error { + if r.hub == nil { + return nil + } + return r.hub.Broadcast(data, conf...) +} + +func (r *Room) GetChannel(channelName string) (*rtmps.Channel, error) { + err := r.LazyInit() + if err != nil { + return nil, err + } + + return r.rtmpa.GetChannel(channelName) } func (r *Room) close() { @@ -72,34 +99,37 @@ func (r *Room) CheckVersion(version uint32) bool { } func (r *Room) UpdateMovie(movieId uint, movie model.BaseMovieInfo) error { + err := r.LazyInit() + if err != nil { + return err + } + m, err := GetMovieByID(r.ID, movieId) if err != nil { return err } switch { case ((m.Live && m.Proxy) || (m.Live && m.RtmpSource)) && (!movie.Live && !movie.Proxy && !movie.RtmpSource): - r.lazyInit() r.rtmpa.DelChannel(m.PullKey) m.PullKey = "" case m.Proxy && !movie.Proxy: m.PullKey = "" } m.MovieInfo.BaseMovieInfo = movie - return db.UpdateMovie(m) + return SaveMovie(m) } -func (r *Room) InitMovie(movie *model.Movie) error { +func (r *Room) initMovie(movie *model.Movie) error { switch { case movie.RtmpSource && movie.Proxy: return errors.New("rtmp source and proxy can't be true at the same time") case movie.Live && movie.RtmpSource: if !conf.Conf.Rtmp.Enable { return errors.New("rtmp is not enabled") - } else if movie.Type == "m3u8" && !conf.Conf.Rtmp.HlsPlayer { - return errors.New("hls player is not enabled") } - movie.PullKey = uuid.New().String() - r.lazyInit() + if movie.PullKey == "" { + movie.PullKey = uuid.New().String() + } _, err := r.rtmpa.NewChannel(movie.PullKey) if err != nil { return err @@ -114,13 +144,13 @@ func (r *Room) InitMovie(movie *model.Movie) error { } switch u.Scheme { case "rtmp": - PullKey := uuid.New().String() - r.lazyInit() - c, err := r.rtmpa.NewChannel(PullKey) + if movie.PullKey == "" { + movie.PullKey = uuid.New().String() + } + c, err := r.rtmpa.NewChannel(movie.PullKey) if err != nil { return err } - movie.PullKey = PullKey go func() { for { if c.Closed() { @@ -139,13 +169,13 @@ func (r *Room) InitMovie(movie *model.Movie) error { } }() case "http", "https": - PullKey := uuid.New().String() - r.lazyInit() - c, err := r.rtmpa.NewChannel(PullKey) + if movie.PullKey == "" { + movie.PullKey = uuid.New().String() + } + c, err := r.rtmpa.NewChannel(movie.PullKey) if err != nil { return err } - movie.PullKey = PullKey go func() { for { if c.Closed() { @@ -176,7 +206,9 @@ func (r *Room) InitMovie(movie *model.Movie) error { if !conf.Conf.Proxy.MovieProxy { return errors.New("movie proxy is not enabled") } - movie.PullKey = uuid.New().String() + if movie.PullKey == "" { + movie.PullKey = uuid.New().String() + } fallthrough case !movie.Live && !movie.Proxy, movie.Live && !movie.Proxy && !movie.RtmpSource: u, err := url.Parse(movie.Url) @@ -193,13 +225,18 @@ func (r *Room) InitMovie(movie *model.Movie) error { } func (r *Room) AddMovie(m model.MovieInfo) error { + err := r.LazyInit() + if err != nil { + return err + } + movie := &model.Movie{ RoomID: r.ID, Position: uint(time.Now().UnixMilli()), MovieInfo: m, } - err := r.InitMovie(movie) + err = r.initMovie(movie) if err != nil { return err } @@ -261,9 +298,9 @@ func (r *Room) GetAllMoviesByRoomID() ([]*model.Movie, error) { if err != nil { return nil, err } - var m []*model.Movie = make([]*model.Movie, ms.Len()) + var m []*model.Movie = make([]*model.Movie, 0, ms.Len()) for i := ms.Front(); i != nil; i = i.Next() { - m[i.Value.Position-1] = i.Value + m = append(m, i.Value) } return m, nil } @@ -277,23 +314,23 @@ func (r *Room) GetMovieByID(id uint) (*model.Movie, error) { } func (r *Room) DeleteMovieByID(id uint) error { + r.LazyInit() m, err := LoadAndDeleteMovieByID(r.ID, id) if err != nil { return err } if m.PullKey != "" { - r.lazyInit() r.rtmpa.DelChannel(m.PullKey) } return nil } func (r *Room) ClearMovies() error { + r.LazyInit() ms, err := db.LoadAndDeleteMoviesByRoomID(r.ID) if err != nil { return err } - r.lazyInit() for _, m := range ms { if m.PullKey != "" { r.rtmpa.DelChannel(m.PullKey) @@ -303,22 +340,22 @@ func (r *Room) ClearMovies() error { } func (r *Room) Current() *Current { - r.lazyInit() c := r.current.Current() return &c } func (r *Room) ChangeCurrentMovie(id uint) error { + r.LazyInit() m, err := GetMovieByID(r.ID, id) if err != nil { return err } - r.lazyInit() r.current.SetMovie(*m) return nil } func (r *Room) SwapMoviePositions(id1, id2 uint) error { + r.LazyInit() return SwapMoviePositions(r.ID, id1, id2) } @@ -327,21 +364,19 @@ func (r *Room) GetMovieWithPullKey(pullKey string) (*model.Movie, error) { } func (r *Room) RegClient(user *User, conn *websocket.Conn) (*Client, error) { - r.lazyInit() + r.LazyInit() return r.hub.RegClient(newClient(user, r, conn)) } func (r *Room) UnregisterClient(user *User) error { - r.lazyInit() + r.LazyInit() return r.hub.UnRegClient(user) } func (r *Room) SetStatus(playing bool, seek float64, rate float64, timeDiff float64) Status { - r.lazyInit() return r.current.SetStatus(playing, seek, rate, timeDiff) } func (r *Room) SetSeekRate(seek float64, rate float64, timeDiff float64) Status { - r.lazyInit() return r.current.SetSeekRate(seek, rate, timeDiff) } diff --git a/internal/op/rooms.go b/internal/op/rooms.go index 8453dbe4..e6f40f83 100644 --- a/internal/op/rooms.go +++ b/internal/op/rooms.go @@ -34,6 +34,7 @@ func initRoom(room *model.Room, conf ...RoomConf) (*Room, error) { Room: *room, lastActive: time.Now().UnixMilli(), version: rand.Uint32(), + current: newCurrent(), } for _, c := range conf { c(r) diff --git a/internal/rtmp/rtmp.go b/internal/rtmp/rtmp.go index 7a2b635b..a7bf999e 100644 --- a/internal/rtmp/rtmp.go +++ b/internal/rtmp/rtmp.go @@ -2,12 +2,10 @@ package rtmp import ( "errors" - "fmt" "strings" "time" "github.com/golang-jwt/jwt/v5" - log "github.com/sirupsen/logrus" "github.com/synctv-org/synctv/internal/conf" rtmps "github.com/zijiren233/livelib/server" "github.com/zijiren233/stream" @@ -46,22 +44,6 @@ func NewRtmpAuthorization(channelName string) (string, error) { func Init(rs *rtmps.Server) { s = rs - - rs.SetParseChannelFunc(func(ReqAppName, ReqChannelName string, IsPublisher bool) (TrueAppName string, TrueChannel string, err error) { - if IsPublisher { - channelName, err := AuthRtmpPublish(ReqChannelName) - if err != nil { - log.Errorf("rtmp: publish auth to %s error: %v", ReqAppName, err) - return "", "", err - } - log.Infof("rtmp: publisher login success: %s/%s", ReqAppName, channelName) - return ReqAppName, channelName, nil - } else if !conf.Conf.Rtmp.RtmpPlayer { - log.Warnf("rtmp: dial to %s/%s error: %s", ReqAppName, ReqChannelName, "rtmp player is not enabled") - return "", "", fmt.Errorf("rtmp: dial to %s/%s error: %s", ReqAppName, ReqChannelName, "rtmp player is not enabled") - } - return ReqAppName, ReqChannelName, nil - }) } func RtmpServer() *rtmps.Server { diff --git a/server/handlers/movie.go b/server/handlers/movie.go index 03e1e6ce..e1ce37aa 100644 --- a/server/handlers/movie.go +++ b/server/handlers/movie.go @@ -149,7 +149,7 @@ func PushMovie(ctx *gin.Context) { return } - if err := room.Hub().Broadcast(&op.ElementMessage{ + if err := room.Broadcast(&op.ElementMessage{ ElementMessage: &pb.ElementMessage{ Type: pb.ElementMessageType_CHANGE_MOVIES, Sender: user.Username, @@ -226,7 +226,7 @@ func EditMovie(ctx *gin.Context) { return } - if err := room.Hub().Broadcast(&op.ElementMessage{ + if err := room.Broadcast(&op.ElementMessage{ ElementMessage: &pb.ElementMessage{ Type: pb.ElementMessageType_CHANGE_MOVIES, Sender: user.Username, @@ -257,7 +257,7 @@ func DelMovie(ctx *gin.Context) { } } - if err := room.Hub().Broadcast(&op.ElementMessage{ + if err := room.Broadcast(&op.ElementMessage{ ElementMessage: &pb.ElementMessage{ Type: pb.ElementMessageType_CHANGE_MOVIES, Sender: user.Username, @@ -279,7 +279,7 @@ func ClearMovies(ctx *gin.Context) { return } - if err := room.Hub().Broadcast(&op.ElementMessage{ + if err := room.Broadcast(&op.ElementMessage{ ElementMessage: &pb.ElementMessage{ Type: pb.ElementMessageType_CHANGE_MOVIES, Sender: user.Username, @@ -307,7 +307,7 @@ func SwapMovie(ctx *gin.Context) { return } - if err := room.Hub().Broadcast(&op.ElementMessage{ + if err := room.Broadcast(&op.ElementMessage{ ElementMessage: &pb.ElementMessage{ Type: pb.ElementMessageType_CHANGE_MOVIES, Sender: user.Username, @@ -334,7 +334,7 @@ func ChangeCurrentMovie(ctx *gin.Context) { ctx.AbortWithStatusJSON(http.StatusBadRequest, model.NewApiErrorResp(err)) return } - if err := room.Hub().Broadcast(&op.ElementMessage{ + if err := room.Broadcast(&op.ElementMessage{ ElementMessage: &pb.ElementMessage{ Type: pb.ElementMessageType_CHANGE_CURRENT, Sender: user.Username, @@ -451,7 +451,7 @@ func JoinLive(ctx *gin.Context) { // ctx.AbortWithStatusJSON(http.StatusNotFound, model.NewApiErrorResp(err)) // return // } - channel, err := room.App().GetChannel(channelName) + channel, err := room.GetChannel(channelName) if err != nil { ctx.AbortWithStatusJSON(http.StatusNotFound, model.NewApiErrorResp(err)) return diff --git a/server/handlers/public.go b/server/handlers/public.go index f23896f6..9b2f737f 100644 --- a/server/handlers/public.go +++ b/server/handlers/public.go @@ -11,7 +11,6 @@ func Settings(ctx *gin.Context) { "rtmp": gin.H{ "enable": conf.Conf.Rtmp.Enable, "rtmpPlayer": conf.Conf.Rtmp.RtmpPlayer, - "hlsPlayer": conf.Conf.Rtmp.HlsPlayer, }, "proxy": gin.H{ "movieProxy": conf.Conf.Proxy.MovieProxy, diff --git a/server/handlers/room.go b/server/handlers/room.go index b7d63ab2..e109508c 100644 --- a/server/handlers/room.go +++ b/server/handlers/room.go @@ -70,7 +70,7 @@ func RoomList(ctx *gin.Context) { resp.Push(&model.RoomListResp{ RoomId: v.ID, RoomName: v.Name, - PeopleNum: v.Hub().ClientNum(), + PeopleNum: v.ClientNum(), NeedPassword: v.NeedPassword(), Creator: op.GetUserName(v.Room.CreatorID), CreatedAt: v.Room.CreatedAt.UnixMilli(), @@ -151,7 +151,7 @@ func CheckRoom(ctx *gin.Context) { } ctx.JSON(http.StatusOK, model.NewApiDataResp(gin.H{ - "peopleNum": r.Hub().ClientNum(), + "peopleNum": r.ClientNum(), "needPassword": r.NeedPassword(), })) } diff --git a/utils/utils.go b/utils/utils.go index 3b89b52e..10575007 100644 --- a/utils/utils.go +++ b/utils/utils.go @@ -176,3 +176,7 @@ func (o *Once) doSlow(f func()) { f() } } + +func (o *Once) Reset() { + atomic.StoreUint32(&o.done, 0) +}