Feat: lazy load movies

pull/21/head
zijiren233 3 years ago
parent 05fbedb2cc
commit 1a1a5f810a

@ -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
}

@ -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
}

@ -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,
}
}

@ -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() {

@ -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)

@ -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)
}

@ -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)

@ -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 {

@ -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

@ -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,

@ -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(),
}))
}

@ -176,3 +176,7 @@ func (o *Once) doSlow(f func()) {
f()
}
}
func (o *Once) Reset() {
atomic.StoreUint32(&o.done, 0)
}

Loading…
Cancel
Save