diff --git a/$null b/$null deleted file mode 100644 index ba3427b..0000000 --- a/$null +++ /dev/null @@ -1,2 +0,0 @@ -'Get-Content' is not recognized as an internal or external command, -operable program or batch file. diff --git a/internal/handler/emby_playback.go b/internal/handler/emby_playback.go index 0450f7a..d9782ba 100644 --- a/internal/handler/emby_playback.go +++ b/internal/handler/emby_playback.go @@ -86,7 +86,9 @@ func embyPrewarmPlaybackTargets(svc *service.Container, c *gin.Context, out map[ if err != nil || m == nil { return } - raw := strings.TrimSpace(m.STRMURL) + // 与播放路径(StreamService.ServeFileWithCloudMode)同一套目标解析: + // STRMURL 为空时回读 .strm 文件内容,预热才能覆盖同一批条目。 + raw := service.MediaSTRMTarget(m) if raw == "" || !service.IsStrmMediaRow(m) { return } @@ -409,9 +411,10 @@ func embyVideoHLSPlaylistHandler(svc *service.Container) gin.HandlerFunc { c.Status(http.StatusNotFound) return } - uid := embyUserID(c) - item, err := svc.Emby.Item(c.Request.Context(), c.Param("id"), uid) - if err != nil || item == nil || svc.Stream == nil { + // 只需确认媒体行存在且对当前用户可见;Emby.Item 会构建完整条目载荷, + // 转码播放下每个分片请求都跑一遍太浪费。 + 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) return } @@ -436,9 +439,8 @@ func embyVideoHLSSegmentHandler(svc *service.Container) gin.HandlerFunc { c.Status(http.StatusNotFound) return } - uid := embyUserID(c) - 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) return } diff --git a/internal/handler/playback_extra.go b/internal/handler/playback_extra.go index 1ecdb2b..0c2c1df 100644 --- a/internal/handler/playback_extra.go +++ b/internal/handler/playback_extra.go @@ -12,6 +12,8 @@ import ( "net/url" "strings" + "go.uber.org/zap" + "github.com/gin-gonic/gin" "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) u, err := svc.Repo.User.FindByID(c.Request.Context(), toString(uid)) 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 "" } token, err := svc.Auth.IssueExternalPlaybackToken(u, mediaID, durationSec) 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 token diff --git a/internal/service/cloud115_hls_proxy.go b/internal/service/cloud115_hls_proxy.go index 626dea1..bce5fba 100644 --- a/internal/service/cloud115_hls_proxy.go +++ b/internal/service/cloud115_hls_proxy.go @@ -128,16 +128,16 @@ func (p *Cloud115HLSProxy) ServeChild(ctx context.Context, w http.ResponseWriter strings.Contains(contentType, "application/vnd.apple.mpegurl") if !isPlaylist { + // Content-Type 缺失时用前 4KB 内容嗅探是否为播放列表;无论结果如何 + // 都要把已读前缀接回 body,避免分片/播放列表丢开头字节。 body, readErr := io.ReadAll(io.LimitReader(resp.Body, 4096)) if readErr != nil { return readErr } if strings.HasPrefix(strings.TrimSpace(string(body)), "#EXTM3U") { 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) @@ -195,9 +195,17 @@ func (p *Cloud115HLSProxy) storeSession(session *cloud115HLSSession) { } } if len(p.sessions) >= cloud115HLSMaxSessions { - for id := range p.sessions { - delete(p.sessions, id) - break + // ExpiresAt 恒为「最近一次活跃时间 + TTL」(lookup 会滑动续期), + // 因此最小者即最久未活跃的会话,淘汰它对正在播放的影响最小。 + 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) @@ -210,9 +218,15 @@ func (p *Cloud115HLSProxy) lookup(sessionID, key string) (*cloud115HLSSession, s } p.mu.Lock() session := p.sessions[sessionID] - if session != nil && time.Now().After(session.ExpiresAt) { - delete(p.sessions, sessionID) - session = nil + if session != nil { + if time.Now().After(session.ExpiresAt) { + delete(p.sessions, sessionID) + session = nil + } else { + // 滑动续期:VOD 点播整个播放期间只有分片请求,master 只在起播时取 + // 一次,不续期的话会话会在 TTL 后被删,长视频播放到一半必然断流。 + session.ExpiresAt = time.Now().Add(cloud115HLSSessionTTL) + } } p.mu.Unlock() if session == nil { diff --git a/internal/service/cloud115_hls_proxy_test.go b/internal/service/cloud115_hls_proxy_test.go index 72225b1..1bfc3ba 100644 --- a/internal/service/cloud115_hls_proxy_test.go +++ b/internal/service/cloud115_hls_proxy_test.go @@ -3,6 +3,7 @@ package service import ( "bytes" "context" + "fmt" "io" "net/http" "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) { for _, host := range []string{"videoplay.115.com", "cpats01.115.com", "cdn.115cdn.net"} { if !isAllowed115UpstreamHost(host) { diff --git a/internal/service/cloud115_playback.go b/internal/service/cloud115_playback.go index 7923f12..af5834e 100644 --- a/internal/service/cloud115_playback.go +++ b/internal/service/cloud115_playback.go @@ -220,6 +220,11 @@ func (s *Cloud115PlaybackService) ResolveCloud115URL(ctx context.Context, mediaI return strings.TrimSpace(item.URL), data, nil } +// cloud115PushStateRetention pushState 条目保留上限。成功条目 3 小时内用于 +// 去重;超过一倍余量即为死数据,触发转码请求时顺带清理,防止 map 随观看过的 +// 条目数无限增长。 +const cloud115PushStateRetention = 6 * time.Hour + func (s *Cloud115PlaybackService) ensureCloudTranscode( ctx context.Context, accountID, pickCode string, @@ -230,6 +235,13 @@ func (s *Cloud115PlaybackService) ensureCloudTranscode( now := time.Now() 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 { ttl := 2 * time.Minute if last.success { diff --git a/internal/service/stream_file.go b/internal/service/stream_file.go index 570c904..a91ea18 100644 --- a/internal/service/stream_file.go +++ b/internal/service/stream_file.go @@ -22,6 +22,25 @@ func (s *StreamService) ServeFile(w http.ResponseWriter, r *http.Request, mediaI 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 { m, err := s.repo.Media.FindByID(r.Context(), mediaID) if err != nil { @@ -30,7 +49,8 @@ func (s *StreamService) ServeFileWithCloudMode(w http.ResponseWriter, r *http.Re if m == nil { 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) { 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) return nil } - if strings.HasPrefix(strings.ToLower(strings.TrimSpace(m.Path)), "cloud://") { - // 云盘媒体没有本地文件可回退;走到这里说明 STRM 播放被关闭或 - // STRMURL 缺失。返回明确错误而不是笼统的「文件不存在」, - // 处理器据此回 502 + 原因,方便用户在播放器/日志里定位。 + pathLower := strings.ToLower(strings.TrimSpace(m.Path)) + if strings.HasPrefix(pathLower, "cloud://") || + strings.HasSuffix(pathLower, ".strm") || + strings.EqualFold(strings.TrimSpace(m.Container), "strm") { + // 云盘/STRM 媒体没有本地视频文件可回退;走到这里说明 STRM 播放被关闭 + // 或播放目标缺失(.strm 内容解析失败)。绝不能把 .strm 文本文件当视频 + // 流出去,返回明确错误而不是笼统的「文件不存在」,处理器据此回 + // 502 + 原因,方便用户在播放器/日志里定位。 return ErrCloudPlaybackUnavailable } f, err := os.Open(m.Path) diff --git a/internal/service/stream_test.go b/internal/service/stream_test.go index 0e20d31..6edbf03 100644 --- a/internal/service/stream_test.go +++ b/internal/service/stream_test.go @@ -5,6 +5,8 @@ import ( "net/http" "net/http/httptest" "net/url" + "os" + "path/filepath" "strings" "testing" @@ -251,6 +253,72 @@ func TestServeFilePassesThroughForeignInstanceSTRMURL(t *testing.T) { // 本机自己生成的 .strm 在换了域名/IP 之后仍要认领:host 对不上,但 acct 是本机 // 网盘账号,于是按当前请求 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) { repos := repository.New(newServiceTestDB(t, &model.Media{}, &model.Setting{}, &model.StrmAccount{})) if err := repos.StrmAccount.Create(t.Context(), &model.StrmAccount{ diff --git a/internal/service/strm_play.go b/internal/service/strm_play.go index 2b3952e..e615514 100644 --- a/internal/service/strm_play.go +++ b/internal/service/strm_play.go @@ -286,7 +286,7 @@ func (s *StrmService) ProxyDirect(ctx context.Context, w http.ResponseWriter, r req.Header.Set("User-Agent", ua) } } - resp, err := s.http.Do(req) + resp, err := s.streamHTTP.Do(req) if err != nil { return err } diff --git a/internal/service/strm_service.go b/internal/service/strm_service.go index 94054f2..9108409 100644 --- a/internal/service/strm_service.go +++ b/internal/service/strm_service.go @@ -11,6 +11,7 @@ import ( "encoding/json" "errors" "fmt" + "net" "net/http" "os" "path/filepath" @@ -81,14 +82,15 @@ var StrmAccountSecretKeys = []string{"cookie", "password", "token", "access_toke // StrmService 提供 STRM 管理的能力。 type StrmService struct { - log *zap.Logger - repo *repository.Container - cfg *config.Config - crypto *CryptoService - http *http.Client - stopOnce sync.Once - stopCh chan struct{} - baseCtx context.Context // 服务级长期上下文(同步/队列不随 HTTP 请求取消) + log *zap.Logger + repo *repository.Container + cfg *config.Config + crypto *CryptoService + http *http.Client + streamHTTP *http.Client // 视频流转发专用:无总超时(见 ProxyDirect) + stopOnce sync.Once + stopCh chan struct{} + baseCtx context.Context // 服务级长期上下文(同步/队列不随 HTTP 请求取消) mu sync.Mutex running map[string]context.CancelFunc // sync path id -> cancel @@ -161,11 +163,24 @@ func (s *StrmService) releaseDownloadSlot(provider string) { // NewStrmService constructs the STRM service. func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService { return &StrmService{ - log: log, - repo: repos, - cfg: cfg, - crypto: crypto, - http: &http.Client{Timeout: 90 * time.Second}, + log: log, + repo: repos, + cfg: cfg, + crypto: crypto, + 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{}), baseCtx: context.Background(), running: map[string]context.CancelFunc{},