From 9c5e3970c6b225678a0542c268734b28da89f682 Mon Sep 17 00:00:00 2001 From: zijiren233 Date: Thu, 24 Oct 2024 00:48:21 +0800 Subject: [PATCH] feat: vendor support proxy m3u8 --- internal/model/movie.go | 4 + server/handlers/movie.go | 94 +------------------ server/handlers/proxy/buffer.go | 59 ++++++++++++ server/handlers/proxy/m3u8.go | 89 ++++++++++++++++++ {utils => server/handlers}/proxy/proxy.go | 35 +------ server/handlers/vendors/vendorAlist/alist.go | 6 +- .../vendors/vendorBilibili/bilibili.go | 2 +- server/handlers/vendors/vendorEmby/emby.go | 4 +- utils/proxy/buffer.go | 24 ----- 9 files changed, 168 insertions(+), 149 deletions(-) create mode 100644 server/handlers/proxy/buffer.go create mode 100644 server/handlers/proxy/m3u8.go rename {utils => server/handlers}/proxy/proxy.go (79%) delete mode 100644 utils/proxy/buffer.go diff --git a/internal/model/movie.go b/internal/model/movie.go index a91dc7f..f77491f 100644 --- a/internal/model/movie.go +++ b/internal/model/movie.go @@ -79,6 +79,10 @@ type MovieBase struct { ParentID EmptyNullString `gorm:"type:char(32)" json:"parentId"` } +func (m *MovieBase) IsM3u8() bool { + return strings.HasPrefix(m.Type, "m3u") || utils.IsM3u8Url(m.Url) +} + func (m *MovieBase) Clone() *MovieBase { mss := make([]*MoreSource, len(m.MoreSources)) for i, ms := range m.MoreSources { diff --git a/server/handlers/movie.go b/server/handlers/movie.go index 83660e2..1d1e8d3 100644 --- a/server/handlers/movie.go +++ b/server/handlers/movie.go @@ -8,30 +8,24 @@ import ( "image" "image/color" "image/png" - "io" "math/rand" "net/http" "path/filepath" "strings" - "time" "github.com/gin-gonic/gin" - "github.com/golang-jwt/jwt/v5" log "github.com/sirupsen/logrus" "github.com/synctv-org/synctv/internal/conf" dbModel "github.com/synctv-org/synctv/internal/model" "github.com/synctv-org/synctv/internal/op" "github.com/synctv-org/synctv/internal/rtmp" "github.com/synctv-org/synctv/internal/settings" + "github.com/synctv-org/synctv/server/handlers/proxy" "github.com/synctv-org/synctv/server/handlers/vendors" "github.com/synctv-org/synctv/server/model" "github.com/synctv-org/synctv/utils" - "github.com/synctv-org/synctv/utils/m3u8" - "github.com/synctv-org/synctv/utils/proxy" - "github.com/zijiren233/go-uhc" "github.com/zijiren233/livelib/protocol/hls" "github.com/zijiren233/livelib/protocol/httpflv" - "github.com/zijiren233/stream" ) func GetPageItems[T any](ctx *gin.Context, items []T) ([]T, error) { @@ -606,14 +600,7 @@ func ProxyMovie(ctx *gin.Context) { // TODO: cache mpd file fallthrough default: - if strings.HasPrefix(m.Movie.MovieBase.Type, "m3u") || utils.IsM3u8Url(m.Movie.MovieBase.Url) { - err = proxyM3u8(ctx, m.Movie.MovieBase.Url, m.Movie.MovieBase.Headers, true, ctx.GetString("token"), room.ID, m.ID) - if err != nil { - log.Errorf("proxy movie error: %v", err) - } - return - } - err = proxy.ProxyURL(ctx, m.Movie.MovieBase.Url, m.Movie.MovieBase.Headers) + err = proxy.AuthProxyURL(ctx, m.Movie.MovieBase.Url, m.Movie.MovieBase.Type, m.Movie.MovieBase.Headers, ctx.GetString("token"), room.ID, m.ID) if err != nil { log.Errorf("proxy movie error: %v", err) return @@ -621,77 +608,6 @@ func ProxyMovie(ctx *gin.Context) { } } -type m3u8TargetClaims struct { - RoomId string `json:"r"` - MovieId string `json:"m"` - TargetUrl string `json:"t"` - jwt.RegisteredClaims -} - -func authM3u8Target(token string) (*m3u8TargetClaims, error) { - t, err := jwt.ParseWithClaims(token, &m3u8TargetClaims{}, func(token *jwt.Token) (any, error) { - return stream.StringToBytes(conf.Conf.Jwt.Secret), nil - }) - if err != nil || !t.Valid { - return nil, ErrAuthFailed - } - claims, ok := t.Claims.(*m3u8TargetClaims) - if !ok { - return nil, ErrAuthFailed - } - return claims, nil -} - -func newM3u8TargetToken(targetUrl, roomId, movieId string) (string, error) { - claims := &m3u8TargetClaims{ - RoomId: roomId, - MovieId: movieId, - TargetUrl: targetUrl, - RegisteredClaims: jwt.RegisteredClaims{ - NotBefore: jwt.NewNumericDate(time.Now()), - }, - } - return jwt.NewWithClaims(jwt.SigningMethodHS256, claims).SignedString(stream.StringToBytes(conf.Conf.Jwt.Secret)) -} - -func proxyM3u8(ctx *gin.Context, u string, headers map[string]string, isM3u8File bool, token, roomId, movieId string) error { - if !isM3u8File { - return proxy.ProxyURL(ctx, u, headers) - } - - req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) - if err != nil { - return fmt.Errorf("new request error: %w", err) - } - for k, v := range headers { - req.Header.Set(k, v) - } - if req.Header.Get("User-Agent") == "" { - req.Header.Set("User-Agent", utils.UA) - } - resp, err := uhc.Do(req) - if err != nil { - return fmt.Errorf("do request error: %w", err) - } - defer resp.Body.Close() - b, err := io.ReadAll(resp.Body) - if err != nil { - return fmt.Errorf("read response body error: %w", err) - } - m3u8Str, err := m3u8.ReplaceM3u8SegmentsWithBaseUrl(stream.BytesToString(b), u, func(segmentUrl string) (string, error) { - targetToken, err := newM3u8TargetToken(segmentUrl, roomId, movieId) - if err != nil { - return "", err - } - return fmt.Sprintf("/api/room/movie/proxy/%s/m3u8/%s?token=%s&roomId=%s", movieId, targetToken, token, roomId), nil - }) - if err != nil { - return fmt.Errorf("replace m3u8 segments with base url error: %w", err) - } - ctx.Data(http.StatusOK, hls.M3U8ContentType, stream.StringToBytes(m3u8Str)) - return nil -} - func ServeM3u8(ctx *gin.Context) { log := ctx.MustGet("log").(*log.Entry) @@ -728,7 +644,7 @@ func ServeM3u8(ctx *gin.Context) { } targetToken := ctx.Param("targetToken") - claims, err := authM3u8Target(targetToken) + claims, err := proxy.GetM3u8Target(targetToken) if err != nil { log.Errorf("auth m3u8 error: %v", err) ctx.AbortWithStatusJSON(http.StatusBadRequest, model.NewApiErrorResp(err)) @@ -738,7 +654,7 @@ func ServeM3u8(ctx *gin.Context) { ctx.AbortWithStatusJSON(http.StatusBadRequest, model.NewApiErrorStringResp("invalid token")) return } - err = proxyM3u8(ctx, claims.TargetUrl, m.Movie.MovieBase.Headers, utils.IsM3u8Url(claims.TargetUrl), ctx.GetString("token"), room.ID, m.ID) + err = proxy.ProxyM3u8(ctx, claims.TargetUrl, m.Movie.MovieBase.Headers, utils.IsM3u8Url(claims.TargetUrl), ctx.GetString("token"), room.ID, m.ID) if err != nil { log.Errorf("proxy m3u8 error: %v", err) } @@ -869,7 +785,7 @@ func JoinHlsLive(ctx *gin.Context) { } if utils.IsM3u8Url(m.Movie.MovieBase.Url) { - _ = proxyM3u8(ctx, m.Movie.MovieBase.Url, m.Movie.MovieBase.Headers, true, ctx.GetString("token"), room.ID, m.ID) + _ = proxy.ProxyM3u8(ctx, m.Movie.MovieBase.Url, m.Movie.MovieBase.Headers, true, ctx.GetString("token"), room.ID, m.ID) return } channel, err := m.Channel() diff --git a/server/handlers/proxy/buffer.go b/server/handlers/proxy/buffer.go new file mode 100644 index 0000000..ee716f4 --- /dev/null +++ b/server/handlers/proxy/buffer.go @@ -0,0 +1,59 @@ +package proxy + +import ( + "errors" + "io" + "sync" +) + +const ( + DefaultBufferSize = 16 * 1024 +) + +var sharedBufferPool = sync.Pool{ + New: func() interface{} { + buffer := make([]byte, DefaultBufferSize) + return &buffer + }, +} + +func getBuffer() *[]byte { + return sharedBufferPool.Get().(*[]byte) +} + +func putBuffer(buffer *[]byte) { + sharedBufferPool.Put(buffer) +} + +func copyBuffer(dst io.Writer, src io.Reader) (written int64, err error) { + buf := getBuffer() + defer putBuffer(buf) + for { + nr, er := src.Read(*buf) + if nr > 0 { + nw, ew := dst.Write((*buf)[0:nr]) + if nw < 0 || nr < nw { + nw = 0 + if ew == nil { + ew = errors.New("invalid write result") + } + } + written += int64(nw) + if ew != nil { + err = ew + break + } + if nr != nw { + err = io.ErrShortWrite + break + } + } + if er != nil { + if er != io.EOF { + err = er + } + break + } + } + return written, err +} diff --git a/server/handlers/proxy/m3u8.go b/server/handlers/proxy/m3u8.go new file mode 100644 index 0000000..19ea503 --- /dev/null +++ b/server/handlers/proxy/m3u8.go @@ -0,0 +1,89 @@ +package proxy + +import ( + "errors" + "fmt" + "io" + "net/http" + "time" + + "github.com/gin-gonic/gin" + "github.com/golang-jwt/jwt/v5" + "github.com/synctv-org/synctv/internal/conf" + "github.com/synctv-org/synctv/utils" + "github.com/synctv-org/synctv/utils/m3u8" + "github.com/zijiren233/go-uhc" + "github.com/zijiren233/livelib/protocol/hls" + "github.com/zijiren233/stream" +) + +type m3u8TargetClaims struct { + RoomId string `json:"r"` + MovieId string `json:"m"` + TargetUrl string `json:"t"` + jwt.RegisteredClaims +} + +func GetM3u8Target(token string) (*m3u8TargetClaims, error) { + t, err := jwt.ParseWithClaims(token, &m3u8TargetClaims{}, func(token *jwt.Token) (any, error) { + return stream.StringToBytes(conf.Conf.Jwt.Secret), nil + }) + if err != nil || !t.Valid { + return nil, errors.New("auth failed") + } + claims, ok := t.Claims.(*m3u8TargetClaims) + if !ok { + return nil, errors.New("auth failed") + } + return claims, nil +} + +func NewM3u8TargetToken(targetUrl, roomId, movieId string) (string, error) { + claims := &m3u8TargetClaims{ + RoomId: roomId, + MovieId: movieId, + TargetUrl: targetUrl, + RegisteredClaims: jwt.RegisteredClaims{ + NotBefore: jwt.NewNumericDate(time.Now()), + }, + } + return jwt.NewWithClaims(jwt.SigningMethodHS256, claims).SignedString(stream.StringToBytes(conf.Conf.Jwt.Secret)) +} + +func ProxyM3u8(ctx *gin.Context, u string, headers map[string]string, isM3u8File bool, token, roomId, movieId string) error { + if !isM3u8File { + return ProxyURL(ctx, u, headers) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) + if err != nil { + return fmt.Errorf("new request error: %w", err) + } + for k, v := range headers { + req.Header.Set(k, v) + } + if req.Header.Get("User-Agent") == "" { + req.Header.Set("User-Agent", utils.UA) + } + resp, err := uhc.Do(req) + if err != nil { + return fmt.Errorf("do request error: %w", err) + } + defer resp.Body.Close() + b, err := io.ReadAll(resp.Body) + if err != nil { + return fmt.Errorf("read response body error: %w", err) + } + m3u8Str, err := m3u8.ReplaceM3u8SegmentsWithBaseUrl(stream.BytesToString(b), u, func(segmentUrl string) (string, error) { + targetToken, err := NewM3u8TargetToken(segmentUrl, roomId, movieId) + if err != nil { + return "", err + } + return fmt.Sprintf("/api/room/movie/proxy/%s/m3u8/%s?token=%s&roomId=%s", movieId, targetToken, token, roomId), nil + }) + if err != nil { + return fmt.Errorf("replace m3u8 segments with base url error: %w", err) + } + ctx.Data(http.StatusOK, hls.M3U8ContentType, stream.StringToBytes(m3u8Str)) + return nil +} diff --git a/utils/proxy/proxy.go b/server/handlers/proxy/proxy.go similarity index 79% rename from utils/proxy/proxy.go rename to server/handlers/proxy/proxy.go index 43eea15..0507033 100644 --- a/utils/proxy/proxy.go +++ b/server/handlers/proxy/proxy.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "net/http" + "strings" "github.com/gin-gonic/gin" "github.com/synctv-org/synctv/internal/settings" @@ -68,35 +69,9 @@ func ProxyURL(ctx *gin.Context, u string, headers map[string]string) error { return nil } -func copyBuffer(dst io.Writer, src io.Reader) (written int64, err error) { - buf := getBuffer() - defer putBuffer(buf) - for { - nr, er := src.Read(*buf) - if nr > 0 { - nw, ew := dst.Write((*buf)[0:nr]) - if nw < 0 || nr < nw { - nw = 0 - if ew == nil { - ew = errors.New("invalid write result") - } - } - written += int64(nw) - if ew != nil { - err = ew - break - } - if nr != nw { - err = io.ErrShortWrite - break - } - } - if er != nil { - if er != io.EOF { - err = er - } - break - } +func AuthProxyURL(ctx *gin.Context, u, t string, headers map[string]string, token, roomId, movieId string) error { + if strings.HasPrefix(t, "m3u") || strings.HasPrefix(u, "m3u") { + return ProxyM3u8(ctx, u, headers, true, token, roomId, movieId) } - return written, err + return ProxyURL(ctx, u, headers) } diff --git a/server/handlers/vendors/vendorAlist/alist.go b/server/handlers/vendors/vendorAlist/alist.go index 1d990e5..2aaddcb 100644 --- a/server/handlers/vendors/vendorAlist/alist.go +++ b/server/handlers/vendors/vendorAlist/alist.go @@ -18,9 +18,9 @@ import ( dbModel "github.com/synctv-org/synctv/internal/model" "github.com/synctv-org/synctv/internal/op" "github.com/synctv-org/synctv/internal/vendor" + "github.com/synctv-org/synctv/server/handlers/proxy" "github.com/synctv-org/synctv/server/model" "github.com/synctv-org/synctv/utils" - "github.com/synctv-org/synctv/utils/proxy" "github.com/synctv-org/vendors/api/alist" ) @@ -134,7 +134,7 @@ func (s *alistVendorService) ProxyMovie(ctx *gin.Context) { ctx.Data(http.StatusOK, "audio/mpegurl", data.Ali.M3U8ListFile) return case "raw": - err := proxy.ProxyURL(ctx, data.URL, nil) + err := proxy.AuthProxyURL(ctx, data.URL, s.movie.MovieBase.Type, nil, ctx.GetString("token"), s.movie.RoomID, s.movie.ID) if err != nil { log.Errorf("proxy vendor movie error: %v", err) } @@ -173,7 +173,7 @@ func (s *alistVendorService) ProxyMovie(ctx *gin.Context) { ctx.AbortWithStatusJSON(http.StatusBadRequest, model.NewApiErrorStringResp("proxy is not enabled")) return } - err = proxy.ProxyURL(ctx, data.URL, nil) + err = proxy.AuthProxyURL(ctx, data.URL, s.movie.MovieBase.Type, nil, ctx.GetString("token"), s.movie.RoomID, s.movie.ID) if err != nil { log.Errorf("proxy vendor movie error: %v", err) } diff --git a/server/handlers/vendors/vendorBilibili/bilibili.go b/server/handlers/vendors/vendorBilibili/bilibili.go index d5cc052..b4700be 100644 --- a/server/handlers/vendors/vendorBilibili/bilibili.go +++ b/server/handlers/vendors/vendorBilibili/bilibili.go @@ -14,9 +14,9 @@ import ( dbModel "github.com/synctv-org/synctv/internal/model" "github.com/synctv-org/synctv/internal/op" "github.com/synctv-org/synctv/internal/vendor" + "github.com/synctv-org/synctv/server/handlers/proxy" "github.com/synctv-org/synctv/server/model" "github.com/synctv-org/synctv/utils" - "github.com/synctv-org/synctv/utils/proxy" "github.com/synctv-org/vendors/api/bilibili" "github.com/zijiren233/stream" "golang.org/x/exp/maps" diff --git a/server/handlers/vendors/vendorEmby/emby.go b/server/handlers/vendors/vendorEmby/emby.go index 3168908..b3c3506 100644 --- a/server/handlers/vendors/vendorEmby/emby.go +++ b/server/handlers/vendors/vendorEmby/emby.go @@ -16,9 +16,9 @@ import ( dbModel "github.com/synctv-org/synctv/internal/model" "github.com/synctv-org/synctv/internal/op" "github.com/synctv-org/synctv/internal/vendor" + "github.com/synctv-org/synctv/server/handlers/proxy" "github.com/synctv-org/synctv/server/model" "github.com/synctv-org/synctv/utils" - "github.com/synctv-org/synctv/utils/proxy" "github.com/synctv-org/vendors/api/emby" ) @@ -151,7 +151,7 @@ func (s *embyVendorService) ProxyMovie(ctx *gin.Context) { ctx.Redirect(http.StatusFound, embyC.Sources[source].URL) return } - err = proxy.ProxyURL(ctx, embyC.Sources[source].URL, nil) + err = proxy.AuthProxyURL(ctx, embyC.Sources[source].URL, "", nil, ctx.GetString("token"), s.movie.RoomID, s.movie.ID) if err != nil { log.Errorf("proxy vendor movie error: %v", err) } diff --git a/utils/proxy/buffer.go b/utils/proxy/buffer.go deleted file mode 100644 index 7d83170..0000000 --- a/utils/proxy/buffer.go +++ /dev/null @@ -1,24 +0,0 @@ -package proxy - -import ( - "sync" -) - -const ( - DefaultBufferSize = 16 * 1024 -) - -var sharedBufferPool = sync.Pool{ - New: func() interface{} { - buffer := make([]byte, DefaultBufferSize) - return &buffer - }, -} - -func getBuffer() *[]byte { - return sharedBufferPool.Get().(*[]byte) -} - -func putBuffer(buffer *[]byte) { - sharedBufferPool.Put(buffer) -}