mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-05 21:06:38 +08:00
优化播放
This commit is contained in:
@@ -1,2 +0,0 @@
|
|||||||
'Get-Content' is not recognized as an internal or external command,
|
|
||||||
operable program or batch file.
|
|
||||||
@@ -86,7 +86,9 @@ func embyPrewarmPlaybackTargets(svc *service.Container, c *gin.Context, out map[
|
|||||||
if err != nil || m == nil {
|
if err != nil || m == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
raw := strings.TrimSpace(m.STRMURL)
|
// 与播放路径(StreamService.ServeFileWithCloudMode)同一套目标解析:
|
||||||
|
// STRMURL 为空时回读 .strm 文件内容,预热才能覆盖同一批条目。
|
||||||
|
raw := service.MediaSTRMTarget(m)
|
||||||
if raw == "" || !service.IsStrmMediaRow(m) {
|
if raw == "" || !service.IsStrmMediaRow(m) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -409,9 +411,10 @@ func embyVideoHLSPlaylistHandler(svc *service.Container) gin.HandlerFunc {
|
|||||||
c.Status(http.StatusNotFound)
|
c.Status(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
uid := embyUserID(c)
|
// 只需确认媒体行存在且对当前用户可见;Emby.Item 会构建完整条目载荷,
|
||||||
item, err := svc.Emby.Item(c.Request.Context(), c.Param("id"), uid)
|
// 转码播放下每个分片请求都跑一遍太浪费。
|
||||||
if err != nil || item == nil || svc.Stream == nil {
|
m, err := svc.Repo.Media.FindByID(c.Request.Context(), c.Param("id"))
|
||||||
|
if err != nil || m == nil || !mediaVisibleForRequest(c, svc, m) || svc.Stream == nil {
|
||||||
c.Status(http.StatusNotFound)
|
c.Status(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -436,9 +439,8 @@ func embyVideoHLSSegmentHandler(svc *service.Container) gin.HandlerFunc {
|
|||||||
c.Status(http.StatusNotFound)
|
c.Status(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
uid := embyUserID(c)
|
m, err := svc.Repo.Media.FindByID(c.Request.Context(), c.Param("id"))
|
||||||
item, err := svc.Emby.Item(c.Request.Context(), c.Param("id"), uid)
|
if err != nil || m == nil || !mediaVisibleForRequest(c, svc, m) || svc.Stream == nil {
|
||||||
if err != nil || item == nil || svc.Stream == nil {
|
|
||||||
c.Status(http.StatusNotFound)
|
c.Status(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,8 @@ import (
|
|||||||
"net/url"
|
"net/url"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
|
"go.uber.org/zap"
|
||||||
|
|
||||||
"github.com/gin-gonic/gin"
|
"github.com/gin-gonic/gin"
|
||||||
|
|
||||||
"github.com/truewhile/MeBox/internal/middleware"
|
"github.com/truewhile/MeBox/internal/middleware"
|
||||||
@@ -173,10 +175,20 @@ func externalPlaybackToken(c *gin.Context, svc *service.Container, mediaID strin
|
|||||||
uid, _ := c.Get(middleware.CtxUserID)
|
uid, _ := c.Get(middleware.CtxUserID)
|
||||||
u, err := svc.Repo.User.FindByID(c.Request.Context(), toString(uid))
|
u, err := svc.Repo.User.FindByID(c.Request.Context(), toString(uid))
|
||||||
if err != nil || u == nil {
|
if err != nil || u == nil {
|
||||||
|
if svc.Log != nil {
|
||||||
|
svc.Log.Warn("external playback token: user lookup failed",
|
||||||
|
zap.String("media_id", mediaID), zap.Error(err))
|
||||||
|
}
|
||||||
return ""
|
return ""
|
||||||
}
|
}
|
||||||
token, err := svc.Auth.IssueExternalPlaybackToken(u, mediaID, durationSec)
|
token, err := svc.Auth.IssueExternalPlaybackToken(u, mediaID, durationSec)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
// 静默返回空串会让 stream_url 变成 ?token=(空)、外部播放器只看到
|
||||||
|
// 401 无从排查,至少在服务端日志里留下原因。
|
||||||
|
if svc.Log != nil {
|
||||||
|
svc.Log.Warn("issue external playback token failed",
|
||||||
|
zap.String("media_id", mediaID), zap.Error(err))
|
||||||
|
}
|
||||||
return ""
|
return ""
|
||||||
}
|
}
|
||||||
return token
|
return token
|
||||||
|
|||||||
@@ -128,16 +128,16 @@ func (p *Cloud115HLSProxy) ServeChild(ctx context.Context, w http.ResponseWriter
|
|||||||
strings.Contains(contentType, "application/vnd.apple.mpegurl")
|
strings.Contains(contentType, "application/vnd.apple.mpegurl")
|
||||||
|
|
||||||
if !isPlaylist {
|
if !isPlaylist {
|
||||||
|
// Content-Type 缺失时用前 4KB 内容嗅探是否为播放列表;无论结果如何
|
||||||
|
// 都要把已读前缀接回 body,避免分片/播放列表丢开头字节。
|
||||||
body, readErr := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
body, readErr := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
||||||
if readErr != nil {
|
if readErr != nil {
|
||||||
return readErr
|
return readErr
|
||||||
}
|
}
|
||||||
if strings.HasPrefix(strings.TrimSpace(string(body)), "#EXTM3U") {
|
if strings.HasPrefix(strings.TrimSpace(string(body)), "#EXTM3U") {
|
||||||
isPlaylist = true
|
isPlaylist = true
|
||||||
resp.Body = io.NopCloser(io.MultiReader(strings.NewReader(string(body)), resp.Body))
|
|
||||||
} else {
|
|
||||||
resp.Body = io.NopCloser(io.MultiReader(strings.NewReader(string(body)), resp.Body))
|
|
||||||
}
|
}
|
||||||
|
resp.Body = io.NopCloser(io.MultiReader(strings.NewReader(string(body)), resp.Body))
|
||||||
}
|
}
|
||||||
|
|
||||||
copyUpstreamHeaders(w, resp, isPlaylist)
|
copyUpstreamHeaders(w, resp, isPlaylist)
|
||||||
@@ -195,9 +195,17 @@ func (p *Cloud115HLSProxy) storeSession(session *cloud115HLSSession) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if len(p.sessions) >= cloud115HLSMaxSessions {
|
if len(p.sessions) >= cloud115HLSMaxSessions {
|
||||||
for id := range p.sessions {
|
// ExpiresAt 恒为「最近一次活跃时间 + TTL」(lookup 会滑动续期),
|
||||||
delete(p.sessions, id)
|
// 因此最小者即最久未活跃的会话,淘汰它对正在播放的影响最小。
|
||||||
break
|
evictID := ""
|
||||||
|
var evictAt time.Time
|
||||||
|
for id, existing := range p.sessions {
|
||||||
|
if evictID == "" || existing.ExpiresAt.Before(evictAt) {
|
||||||
|
evictID, evictAt = id, existing.ExpiresAt
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if evictID != "" {
|
||||||
|
delete(p.sessions, evictID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
session.ExpiresAt = now.Add(cloud115HLSSessionTTL)
|
session.ExpiresAt = now.Add(cloud115HLSSessionTTL)
|
||||||
@@ -210,9 +218,15 @@ func (p *Cloud115HLSProxy) lookup(sessionID, key string) (*cloud115HLSSession, s
|
|||||||
}
|
}
|
||||||
p.mu.Lock()
|
p.mu.Lock()
|
||||||
session := p.sessions[sessionID]
|
session := p.sessions[sessionID]
|
||||||
if session != nil && time.Now().After(session.ExpiresAt) {
|
if session != nil {
|
||||||
delete(p.sessions, sessionID)
|
if time.Now().After(session.ExpiresAt) {
|
||||||
session = nil
|
delete(p.sessions, sessionID)
|
||||||
|
session = nil
|
||||||
|
} else {
|
||||||
|
// 滑动续期:VOD 点播整个播放期间只有分片请求,master 只在起播时取
|
||||||
|
// 一次,不续期的话会话会在 TTL 后被删,长视频播放到一半必然断流。
|
||||||
|
session.ExpiresAt = time.Now().Add(cloud115HLSSessionTTL)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
p.mu.Unlock()
|
p.mu.Unlock()
|
||||||
if session == nil {
|
if session == nil {
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package service
|
|||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
@@ -56,6 +57,67 @@ https://cpats01.115.com/seg.ts?x=1
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 长视频 VOD 播放期间只有分片请求经过 lookup,会话必须随活动滑动续期,
|
||||||
|
// 否则固定 30 分钟 TTL 一到就会在中途 404 断流。
|
||||||
|
func TestCloud115HLSProxyLookupRenewsSession(t *testing.T) {
|
||||||
|
proxy := &Cloud115HLSProxy{sessions: map[string]*cloud115HLSSession{}}
|
||||||
|
session := &cloud115HLSSession{
|
||||||
|
ID: "sess",
|
||||||
|
MediaID: "media-1",
|
||||||
|
entries: map[string]string{"e1": "https://cpats01.115.com/v.m3u8"},
|
||||||
|
}
|
||||||
|
proxy.storeSession(session)
|
||||||
|
// 模拟播放进行到第 29 分钟:距过期只剩 1 分钟。
|
||||||
|
session.ExpiresAt = time.Now().Add(time.Minute)
|
||||||
|
|
||||||
|
if _, _, ok := proxy.lookup(session.ID, "e1"); !ok {
|
||||||
|
t.Fatal("active session should still resolve before expiry")
|
||||||
|
}
|
||||||
|
if remaining := time.Until(session.ExpiresAt); remaining < cloud115HLSSessionTTL-10*time.Second {
|
||||||
|
t.Fatalf("session not renewed on activity, remaining = %v", remaining)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 真正过期(期间无任何活动)的会话仍要被清理。
|
||||||
|
session.ExpiresAt = time.Now().Add(-time.Minute)
|
||||||
|
if _, _, ok := proxy.lookup(session.ID, "e1"); ok {
|
||||||
|
t.Fatal("expired session should not resolve")
|
||||||
|
}
|
||||||
|
proxy.mu.Lock()
|
||||||
|
_, exists := proxy.sessions[session.ID]
|
||||||
|
proxy.mu.Unlock()
|
||||||
|
if exists {
|
||||||
|
t.Fatal("expired session should be evicted from the map")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 满员驱逐必须淘汰最久未活跃的会话(ExpiresAt = 最近活跃 + TTL,最小者即
|
||||||
|
// 最久未活跃),而不是随机挑一个——随机驱逐可能正好杀掉正在播放的会话。
|
||||||
|
func TestCloud115HLSProxyEvictsLeastRecentlyActiveSession(t *testing.T) {
|
||||||
|
proxy := &Cloud115HLSProxy{sessions: map[string]*cloud115HLSSession{}}
|
||||||
|
first := &cloud115HLSSession{ID: "active", MediaID: "media-1", entries: map[string]string{}}
|
||||||
|
proxy.storeSession(first)
|
||||||
|
for i := 0; i < cloud115HLSMaxSessions-1; i++ {
|
||||||
|
proxy.storeSession(&cloud115HLSSession{ID: fmt.Sprintf("s-%03d", i), MediaID: "media-1", entries: map[string]string{}})
|
||||||
|
}
|
||||||
|
// 人为拉开活跃时间差,避免依赖时钟精度:active 最近活跃,s-000 最久未活跃。
|
||||||
|
base := time.Now()
|
||||||
|
first.ExpiresAt = base.Add(cloud115HLSSessionTTL)
|
||||||
|
proxy.sessions["s-000"].ExpiresAt = base.Add(cloud115HLSSessionTTL - time.Hour)
|
||||||
|
|
||||||
|
proxy.storeSession(&cloud115HLSSession{ID: "newcomer", MediaID: "media-1", entries: map[string]string{}})
|
||||||
|
|
||||||
|
proxy.mu.Lock()
|
||||||
|
_, activeExists := proxy.sessions["active"]
|
||||||
|
_, staleExists := proxy.sessions["s-000"]
|
||||||
|
proxy.mu.Unlock()
|
||||||
|
if !activeExists {
|
||||||
|
t.Fatal("recently active session must survive eviction")
|
||||||
|
}
|
||||||
|
if staleExists {
|
||||||
|
t.Fatal("least recently active session should have been evicted")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestIsAllowed115UpstreamHost(t *testing.T) {
|
func TestIsAllowed115UpstreamHost(t *testing.T) {
|
||||||
for _, host := range []string{"videoplay.115.com", "cpats01.115.com", "cdn.115cdn.net"} {
|
for _, host := range []string{"videoplay.115.com", "cpats01.115.com", "cdn.115cdn.net"} {
|
||||||
if !isAllowed115UpstreamHost(host) {
|
if !isAllowed115UpstreamHost(host) {
|
||||||
|
|||||||
@@ -220,6 +220,11 @@ func (s *Cloud115PlaybackService) ResolveCloud115URL(ctx context.Context, mediaI
|
|||||||
return strings.TrimSpace(item.URL), data, nil
|
return strings.TrimSpace(item.URL), data, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// cloud115PushStateRetention pushState 条目保留上限。成功条目 3 小时内用于
|
||||||
|
// 去重;超过一倍余量即为死数据,触发转码请求时顺带清理,防止 map 随观看过的
|
||||||
|
// 条目数无限增长。
|
||||||
|
const cloud115PushStateRetention = 6 * time.Hour
|
||||||
|
|
||||||
func (s *Cloud115PlaybackService) ensureCloudTranscode(
|
func (s *Cloud115PlaybackService) ensureCloudTranscode(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
accountID, pickCode string,
|
accountID, pickCode string,
|
||||||
@@ -230,6 +235,13 @@ func (s *Cloud115PlaybackService) ensureCloudTranscode(
|
|||||||
now := time.Now()
|
now := time.Now()
|
||||||
|
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
|
if len(s.pushState) > 0 {
|
||||||
|
for k, v := range s.pushState {
|
||||||
|
if now.Sub(v.at) > cloud115PushStateRetention {
|
||||||
|
delete(s.pushState, k)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
if last, ok := s.pushState[key]; ok {
|
if last, ok := s.pushState[key]; ok {
|
||||||
ttl := 2 * time.Minute
|
ttl := 2 * time.Minute
|
||||||
if last.success {
|
if last.success {
|
||||||
|
|||||||
@@ -22,6 +22,25 @@ func (s *StreamService) ServeFile(w http.ResponseWriter, r *http.Request, mediaI
|
|||||||
return s.ServeFileWithCloudMode(w, r, mediaID, "")
|
return s.ServeFileWithCloudMode(w, r, mediaID, "")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MediaSTRMTarget 返回媒体行固化的 STRM 播放目标。STRMURL 为空时回读本地
|
||||||
|
// .strm 文件内容兜底(扫描时内容解析失败的行只剩 Container=strm + Path),
|
||||||
|
// 仍拿不到返回空串。
|
||||||
|
func MediaSTRMTarget(m *model.Media) string {
|
||||||
|
if m == nil {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
if raw := strings.TrimSpace(m.STRMURL); raw != "" {
|
||||||
|
return raw
|
||||||
|
}
|
||||||
|
path := strings.TrimSpace(m.Path)
|
||||||
|
if strings.HasSuffix(strings.ToLower(path), ".strm") {
|
||||||
|
if target, err := readLocalSTRMTarget(path); err == nil {
|
||||||
|
return strings.TrimSpace(target)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Request, mediaID, cloudMode string) error {
|
func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Request, mediaID, cloudMode string) error {
|
||||||
m, err := s.repo.Media.FindByID(r.Context(), mediaID)
|
m, err := s.repo.Media.FindByID(r.Context(), mediaID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -30,7 +49,8 @@ func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Re
|
|||||||
if m == nil {
|
if m == nil {
|
||||||
return ErrMediaNotFound
|
return ErrMediaNotFound
|
||||||
}
|
}
|
||||||
if strmURL := strings.TrimSpace(m.STRMURL); strmURL != "" && playableSTRMTarget(r.Context(), s.repo, strmURL, m) {
|
strmURL := MediaSTRMTarget(m)
|
||||||
|
if strmURL != "" && playableSTRMTarget(r.Context(), s.repo, strmURL, m) {
|
||||||
if !cloudPlaybackModeEnabled(r.Context(), s.repo, cloudMode) {
|
if !cloudPlaybackModeEnabled(r.Context(), s.repo, cloudMode) {
|
||||||
return ErrCloudPlaybackDisabled
|
return ErrCloudPlaybackDisabled
|
||||||
}
|
}
|
||||||
@@ -50,10 +70,14 @@ func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Re
|
|||||||
http.Redirect(w, r, absoluteInternalRedirect(target, r), http.StatusFound)
|
http.Redirect(w, r, absoluteInternalRedirect(target, r), http.StatusFound)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
if strings.HasPrefix(strings.ToLower(strings.TrimSpace(m.Path)), "cloud://") {
|
pathLower := strings.ToLower(strings.TrimSpace(m.Path))
|
||||||
// 云盘媒体没有本地文件可回退;走到这里说明 STRM 播放被关闭或
|
if strings.HasPrefix(pathLower, "cloud://") ||
|
||||||
// STRMURL 缺失。返回明确错误而不是笼统的「文件不存在」,
|
strings.HasSuffix(pathLower, ".strm") ||
|
||||||
// 处理器据此回 502 + 原因,方便用户在播放器/日志里定位。
|
strings.EqualFold(strings.TrimSpace(m.Container), "strm") {
|
||||||
|
// 云盘/STRM 媒体没有本地视频文件可回退;走到这里说明 STRM 播放被关闭
|
||||||
|
// 或播放目标缺失(.strm 内容解析失败)。绝不能把 .strm 文本文件当视频
|
||||||
|
// 流出去,返回明确错误而不是笼统的「文件不存在」,处理器据此回
|
||||||
|
// 502 + 原因,方便用户在播放器/日志里定位。
|
||||||
return ErrCloudPlaybackUnavailable
|
return ErrCloudPlaybackUnavailable
|
||||||
}
|
}
|
||||||
f, err := os.Open(m.Path)
|
f, err := os.Open(m.Path)
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
@@ -251,6 +253,72 @@ func TestServeFilePassesThroughForeignInstanceSTRMURL(t *testing.T) {
|
|||||||
|
|
||||||
// 本机自己生成的 .strm 在换了域名/IP 之后仍要认领:host 对不上,但 acct 是本机
|
// 本机自己生成的 .strm 在换了域名/IP 之后仍要认领:host 对不上,但 acct 是本机
|
||||||
// 网盘账号,于是按当前请求 host 相对化,保持可播放。
|
// 网盘账号,于是按当前请求 host 相对化,保持可播放。
|
||||||
|
// 扫描时 .strm 内容解析失败会产生 Container=strm 但 STRMURL 为空的行:
|
||||||
|
// 播放路径必须回读 .strm 文件内容兜底(与 MediaPlaybackProvider 等一致),
|
||||||
|
// 而不是把 .strm 文本文件当视频流返回。
|
||||||
|
func TestServeFileFallsBackToStrmFileContentWhenSTRMURLEmpty(t *testing.T) {
|
||||||
|
repos := newStreamTestRepo(t)
|
||||||
|
dir := t.TempDir()
|
||||||
|
strmPath := filepath.Join(dir, "Movie.strm")
|
||||||
|
target := "https://cdn.example.test/Movie.mkv?sign=direct"
|
||||||
|
if err := os.WriteFile(strmPath, []byte(target+"\n"), 0o600); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := repos.DB.Create(&model.Media{
|
||||||
|
Base: model.Base{ID: "strm-empty-url"},
|
||||||
|
Title: "STRM Empty URL",
|
||||||
|
Path: strmPath,
|
||||||
|
Container: "strm",
|
||||||
|
STRMURL: "",
|
||||||
|
}).Error; err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
svc := NewStreamService(&config.Config{}, zap.NewNop(), repos, nil)
|
||||||
|
req := httptest.NewRequest(http.MethodGet, "http://nas.local:18080/api/stream/strm-empty-url?token=jwt123", nil)
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
|
||||||
|
if err := svc.ServeFile(w, req, "strm-empty-url"); err != nil {
|
||||||
|
t.Fatalf("empty-STRMURL .strm row should fall back to file content: %v", err)
|
||||||
|
}
|
||||||
|
if w.Code != http.StatusFound {
|
||||||
|
t.Fatalf("status = %d, want 302", w.Code)
|
||||||
|
}
|
||||||
|
if loc := w.Header().Get("Location"); loc != target {
|
||||||
|
t.Fatalf("Location = %q, want %q", loc, target)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// .strm 内容完全解析不出播放目标时,返回明确的 502 错误,绝不把 .strm 文本
|
||||||
|
// 文件本身当视频流吐给播放器。
|
||||||
|
func TestServeFileRejectsStrmRowWithUnresolvableTarget(t *testing.T) {
|
||||||
|
repos := newStreamTestRepo(t)
|
||||||
|
dir := t.TempDir()
|
||||||
|
strmPath := filepath.Join(dir, "Broken.strm")
|
||||||
|
if err := os.WriteFile(strmPath, []byte("# 只有注释,没有可用播放地址\n"), 0o600); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := repos.DB.Create(&model.Media{
|
||||||
|
Base: model.Base{ID: "strm-broken"},
|
||||||
|
Title: "STRM Broken",
|
||||||
|
Path: strmPath,
|
||||||
|
Container: "strm",
|
||||||
|
STRMURL: "",
|
||||||
|
}).Error; err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
svc := NewStreamService(&config.Config{}, zap.NewNop(), repos, nil)
|
||||||
|
req := httptest.NewRequest(http.MethodGet, "http://nas.local:18080/api/stream/strm-broken", nil)
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
|
||||||
|
err := svc.ServeFile(w, req, "strm-broken")
|
||||||
|
if !errors.Is(err, ErrCloudPlaybackUnavailable) {
|
||||||
|
t.Fatalf("error = %v, want ErrCloudPlaybackUnavailable", err)
|
||||||
|
}
|
||||||
|
if w.Code != http.StatusOK {
|
||||||
|
t.Fatalf("no bytes should be written on error, status = %d", w.Code)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestServeFileRealignsOwnSTRMURLOtherHost(t *testing.T) {
|
func TestServeFileRealignsOwnSTRMURLOtherHost(t *testing.T) {
|
||||||
repos := repository.New(newServiceTestDB(t, &model.Media{}, &model.Setting{}, &model.StrmAccount{}))
|
repos := repository.New(newServiceTestDB(t, &model.Media{}, &model.Setting{}, &model.StrmAccount{}))
|
||||||
if err := repos.StrmAccount.Create(t.Context(), &model.StrmAccount{
|
if err := repos.StrmAccount.Create(t.Context(), &model.StrmAccount{
|
||||||
|
|||||||
@@ -286,7 +286,7 @@ func (s *StrmService) ProxyDirect(ctx context.Context, w http.ResponseWriter, r
|
|||||||
req.Header.Set("User-Agent", ua)
|
req.Header.Set("User-Agent", ua)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
resp, err := s.http.Do(req)
|
resp, err := s.streamHTTP.Do(req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
@@ -81,14 +82,15 @@ var StrmAccountSecretKeys = []string{"cookie", "password", "token", "access_toke
|
|||||||
|
|
||||||
// StrmService 提供 STRM 管理的能力。
|
// StrmService 提供 STRM 管理的能力。
|
||||||
type StrmService struct {
|
type StrmService struct {
|
||||||
log *zap.Logger
|
log *zap.Logger
|
||||||
repo *repository.Container
|
repo *repository.Container
|
||||||
cfg *config.Config
|
cfg *config.Config
|
||||||
crypto *CryptoService
|
crypto *CryptoService
|
||||||
http *http.Client
|
http *http.Client
|
||||||
stopOnce sync.Once
|
streamHTTP *http.Client // 视频流转发专用:无总超时(见 ProxyDirect)
|
||||||
stopCh chan struct{}
|
stopOnce sync.Once
|
||||||
baseCtx context.Context // 服务级长期上下文(同步/队列不随 HTTP 请求取消)
|
stopCh chan struct{}
|
||||||
|
baseCtx context.Context // 服务级长期上下文(同步/队列不随 HTTP 请求取消)
|
||||||
|
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
running map[string]context.CancelFunc // sync path id -> cancel
|
running map[string]context.CancelFunc // sync path id -> cancel
|
||||||
@@ -161,11 +163,24 @@ func (s *StrmService) releaseDownloadSlot(provider string) {
|
|||||||
// NewStrmService constructs the STRM service.
|
// NewStrmService constructs the STRM service.
|
||||||
func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService {
|
func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService {
|
||||||
return &StrmService{
|
return &StrmService{
|
||||||
log: log,
|
log: log,
|
||||||
repo: repos,
|
repo: repos,
|
||||||
cfg: cfg,
|
cfg: cfg,
|
||||||
crypto: crypto,
|
crypto: crypto,
|
||||||
http: &http.Client{Timeout: 90 * time.Second},
|
http: &http.Client{Timeout: 90 * time.Second},
|
||||||
|
// 视频流转发不能用带总超时的 client:http.Client.Timeout 覆盖整个
|
||||||
|
// 响应体读取,长视频必然超过 90s 被硬切。传输生命周期由请求 ctx
|
||||||
|
// (客户端断开即取消)控制,这里只保留建连/响应头阶段的兜底超时。
|
||||||
|
streamHTTP: &http.Client{
|
||||||
|
Transport: &http.Transport{
|
||||||
|
Proxy: http.ProxyFromEnvironment,
|
||||||
|
DialContext: (&net.Dialer{Timeout: 15 * time.Second, KeepAlive: 30 * time.Second}).DialContext,
|
||||||
|
ForceAttemptHTTP2: true,
|
||||||
|
TLSHandshakeTimeout: 10 * time.Second,
|
||||||
|
ResponseHeaderTimeout: 30 * time.Second,
|
||||||
|
IdleConnTimeout: 90 * time.Second,
|
||||||
|
},
|
||||||
|
},
|
||||||
stopCh: make(chan struct{}),
|
stopCh: make(chan struct{}),
|
||||||
baseCtx: context.Background(),
|
baseCtx: context.Background(),
|
||||||
running: map[string]context.CancelFunc{},
|
running: map[string]context.CancelFunc{},
|
||||||
|
|||||||
Reference in New Issue
Block a user