perf(strm): 优化 OpenList 元数据下载速度 (#15)

* Add favourite buttons to library list and series detail pages

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Fix favourites list for remote Emby mounted media

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Route favourites to library series detail when appropriate

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Fix favourites link to library series for remote anime paths

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Fix remote Emby playback history and continue watching

- Preserve duration_ms when web player reports duration=0
- Add GET /playback/:id/resume for reliable per-item resume
- Switch home continue watching to /watch-history/continue API
- Set library_id on remote media for play-profile visibility
- Keep continue-watching rows when remote metadata hydration fails

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Fix mobile sort dropdown overflowing off-screen

Align dropdown to button left edge on small screens and cap width
to viewport so sort options stay fully visible on mobile.

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Remove stale Windows binary backup files

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* Fix STRM account edit form not echoing saved config

Return config_preview from account API for server/url/username and
secret flags. Merge partial config updates so empty password/token
fields do not wipe stored credentials.

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

* perf(strm): speed up OpenList metadata downloads with provider-aware concurrency

- Split download semaphore: 115 stays capped at 3 for WAF safety; OpenList/CloudDrive2
  WebDAV metadata uses strm.download_threads (default raised to 6, max 16)
- Scope 115 WAF cooldown to 115 tasks only so OpenList downloads are not paused
- Requeue claimed tasks back to pending when slot acquisition fails or 115 is cooling
- Add tests for separate semaphores and requeue behavior

Co-authored-by: truewhile <truewhile@users.noreply.github.com>

---------

Co-authored-by: Cursor Agent <cursoragent@cursor.com>
Co-authored-by: truewhile <truewhile@users.noreply.github.com>
This commit is contained in:
truewhile
2026-09-02 15:15:49 +08:00
committed by GitHub
parent fbc33862fb
commit 67b85840dc
23 changed files with 515 additions and 91 deletions
+24
View File
@@ -82,6 +82,30 @@ func playbackProgressHandler(svc *service.Container) gin.HandlerFunc {
}
}
func playbackResumeHandler(svc *service.Container) gin.HandlerFunc {
return func(c *gin.Context) {
uid, _ := c.Get(middleware.CtxUserID)
row, err := svc.Playback.GetProgress(c.Request.Context(), toString(uid), c.Param("id"))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
if row == nil {
c.JSON(http.StatusOK, gin.H{
"position_ms": 0,
"duration_ms": 0,
"completed": false,
})
return
}
c.JSON(http.StatusOK, gin.H{
"position_ms": row.PositionMs,
"duration_ms": row.DurationMs,
"completed": row.Completed,
})
}
}
// externalPlayersHandler returns the list of external player URI
// schemes the UI can offer the user. We lookup the media row to
// produce the per-player launch URL.
@@ -71,6 +71,7 @@ func registerAuthedFavoriteAndMediaActionRoutes(authed *gin.RouterGroup, svc *se
func registerAuthedPlaybackExtraRoutes(authed *gin.RouterGroup, svc *service.Container) {
authed.GET("/playback/:id/info", playbackInfoHandler(svc))
authed.GET("/playback/:id/resume", playbackResumeHandler(svc))
authed.POST("/playback/:id/progress", playbackProgressHandler(svc))
authed.GET("/playback/:id/external-players", externalPlayersHandler(svc))
authed.GET("/playback/:id/external-url", externalURLHandler(svc))
+5
View File
@@ -29,6 +29,8 @@ type strmAccountView struct {
model.StrmAccount
HasCredential bool `json:"has_credential"`
ProviderLabel string `json:"provider_label"`
// ConfigPreview 非敏感配置字段,供编辑表单回显。
ConfigPreview service.StrmAccountConfigPreview `json:"config_preview,omitempty"`
// ProxyPlay 仅远程 Emby 挂载账号返回:播放流量是否经过 MMTL 代理(编辑回显用)。
ProxyPlay *bool `json:"proxy_play,omitempty"`
// EmbyLines 仅远程 Emby 挂载账号返回:多线路配置(不含凭据)。
@@ -45,6 +47,9 @@ func strmAccountViews(svc *service.Container, accounts []model.StrmAccount) []st
HasCredential: service.HasStrmAccountCredential(&a),
ProviderLabel: providerLabelOf(a.Provider),
}
if svc != nil && svc.Strm != nil {
view.ConfigPreview = svc.Strm.StrmAccountConfigPreviewOf(&a)
}
if a.Provider == model.StrmProviderEmbyRemote && svc != nil && svc.EmbyRemote != nil {
if proxyPlay, err := svc.EmbyRemote.ProxyPlayOf(&a); err == nil {
view.ProxyPlay = &proxyPlay
+27 -4
View File
@@ -129,14 +129,23 @@ func historyContinueHandler(svc *service.Container) gin.HandlerFunc {
mountID, remoteID, _ := service.DecodeEmbyRemoteID(r.MediaID)
if mount, acct, _ := svc.EmbyRemote.ResolveMount(c.Request.Context(), mountID); mount != nil && acct != nil {
if rm, err := svc.EmbyRemote.RemoteMediaDetail(c.Request.Context(), mount, acct, remoteID); err == nil && rm != nil {
out = append(out, gin.H{
"history": r,
"media": *rm,
})
if mediaVisibleForRequest(c, svc, rm) {
out = append(out, gin.H{
"history": r,
"media": *rm,
})
}
continue
}
}
}
fallback := fallbackHistoryMedia(r.MediaID)
if fallback != nil {
out = append(out, gin.H{
"history": r,
"media": *fallback,
})
}
continue
}
out = append(out, gin.H{
@@ -148,6 +157,20 @@ func historyContinueHandler(svc *service.Container) gin.HandlerFunc {
}
}
func fallbackHistoryMedia(mediaID string) *model.Media {
if mediaID == "" {
return nil
}
title := "媒体"
if service.IsEmbyRemoteID(mediaID) {
title = "远程媒体"
}
return &model.Media{
Base: model.Base{ID: mediaID},
Title: title,
}
}
// historyDeleteHandler removes one or all history rows for the caller.
//
// DELETE /api/watch-history?media_id=xxx → delete just that media's row
+3 -1
View File
@@ -26,7 +26,9 @@ func (r *HistoryRepository) Upsert(ctx context.Context, h *model.PlaybackHistory
return err
}
existing.PositionMs = h.PositionMs
existing.DurationMs = h.DurationMs
if h.DurationMs > 0 {
existing.DurationMs = h.DurationMs
}
existing.WatchedAt = h.WatchedAt
existing.Completed = h.Completed
return r.db.WithContext(ctx).Save(&existing).Error
+3 -1
View File
@@ -301,7 +301,9 @@ func (r *EmbyRemoteService) MapRemoteItemToMedia(ctx context.Context, mount *mod
media.EpisodeNum = 0
}
if mount != nil && strings.TrimSpace(mount.RemoteViewID) != "" {
media.DisplayLibraryID = EncodeEmbyRemoteID(mount.ID, mount.RemoteViewID)
libID := EncodeEmbyRemoteID(mount.ID, mount.RemoteViewID)
media.DisplayLibraryID = libID
media.LibraryID = libID
}
return media
}
+46 -2
View File
@@ -13,6 +13,7 @@ import (
"time"
"go.uber.org/zap"
"gorm.io/gorm"
"github.com/ShukeBta/MMTL/internal/model"
"github.com/ShukeBta/MMTL/internal/repository"
@@ -47,18 +48,61 @@ func (p *PlaybackService) RecordProgress(ctx context.Context, userID, mediaID st
if userID == "" || mediaID == "" {
return errors.New("missing user or media")
}
completed := duration > 0 && position >= duration-30_000
dur := p.resolvePlaybackDuration(ctx, userID, mediaID, duration)
completed := dur > 0 && position >= dur-30_000
h := &model.PlaybackHistory{
UserID: userID,
MediaID: mediaID,
PositionMs: position,
DurationMs: duration,
DurationMs: dur,
WatchedAt: time.Now(),
Completed: completed,
}
return p.repo.History.Upsert(ctx, h)
}
// GetProgress returns the saved resume row for one media item, or nil when absent.
func (p *PlaybackService) GetProgress(ctx context.Context, userID, mediaID string) (*model.PlaybackHistory, error) {
if userID == "" || mediaID == "" {
return nil, errors.New("missing user or media")
}
var row model.PlaybackHistory
err := p.repo.DB.WithContext(ctx).
Where("user_id = ? AND media_id = ?", userID, mediaID).
First(&row).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil, nil
}
if err != nil {
return nil, err
}
return &row, nil
}
func (p *PlaybackService) resolvePlaybackDuration(ctx context.Context, userID, mediaID string, duration int64) int64 {
if duration > 0 {
return duration
}
var existing model.PlaybackHistory
if err := p.repo.DB.WithContext(ctx).
Where("user_id = ? AND media_id = ?", userID, mediaID).
First(&existing).Error; err == nil && existing.DurationMs > 0 {
return existing.DurationMs
}
if m, _ := p.repo.Media.FindByID(ctx, mediaID); m != nil && m.DurationSec > 0 {
return int64(m.DurationSec) * 1000
}
if p.remote != nil && IsEmbyRemoteID(mediaID) {
mountID, remoteID, _ := DecodeEmbyRemoteID(mediaID)
if mount, acct, _ := p.remote.ResolveMount(ctx, mountID); mount != nil && acct != nil {
if rm, err := p.remote.RemoteMediaDetail(ctx, mount, acct, remoteID); err == nil && rm != nil && rm.DurationSec > 0 {
return int64(rm.DurationSec) * 1000
}
}
}
return duration
}
// HistoryItem joins the playback row with its media so the API consumer
// gets a fully-populated card without a second round-trip.
type HistoryItem struct {
+47
View File
@@ -10,6 +10,53 @@ import (
"go.uber.org/zap"
)
func TestRecordProgressPreservesDurationWhenMissing(t *testing.T) {
db := newServiceTestDB(t, &model.PlaybackHistory{}, &model.Media{})
repos := repository.New(db)
userID := "user-1"
mediaID := "local-movie-1"
if err := db.Create(&model.PlaybackHistory{
Base: model.Base{ID: "hist-1"},
UserID: userID,
MediaID: mediaID,
PositionMs: 30_000,
DurationMs: 120_000,
WatchedAt: time.Now(),
}).Error; err != nil {
t.Fatalf("seed history: %v", err)
}
svc := NewPlaybackService(zap.NewNop(), repos)
if err := svc.RecordProgress(context.Background(), userID, mediaID, 60_000, 0); err != nil {
t.Fatalf("RecordProgress: %v", err)
}
var hist model.PlaybackHistory
if err := db.Where("user_id = ? AND media_id = ?", userID, mediaID).First(&hist).Error; err != nil {
t.Fatalf("find history: %v", err)
}
if hist.DurationMs != 120_000 {
t.Fatalf("expected duration preserved, got %d", hist.DurationMs)
}
if hist.PositionMs != 60_000 {
t.Fatalf("expected position updated, got %d", hist.PositionMs)
}
}
func TestGetProgressReturnsNilWhenMissing(t *testing.T) {
db := newServiceTestDB(t, &model.PlaybackHistory{})
repos := repository.New(db)
svc := NewPlaybackService(zap.NewNop(), repos)
row, err := svc.GetProgress(context.Background(), "user-1", "missing-media")
if err != nil {
t.Fatalf("GetProgress: %v", err)
}
if row != nil {
t.Fatalf("expected nil progress, got %#v", row)
}
}
func TestListFavouritesIncludesLocalAndRemoteIDs(t *testing.T) {
db := newServiceTestDB(t, &model.Favorite{}, &model.Media{})
repos := repository.New(db)
+76
View File
@@ -0,0 +1,76 @@
package service
import (
"context"
"encoding/json"
"testing"
"github.com/ShukeBta/MMTL/internal/model"
)
func TestStrmAccountConfigPreviewOf(t *testing.T) {
svc := testStrmService(t)
acct := &model.StrmAccount{
Provider: model.StrmProviderOpenList,
Config: mustJSON(map[string]string{
"server": "https://list.example.com",
"username": "alice",
"password": svc.crypto.Encrypt("secret"),
"token": svc.crypto.Encrypt("tok"),
}),
}
preview := svc.StrmAccountConfigPreviewOf(acct)
if preview.Server != "https://list.example.com" {
t.Fatalf("server = %q", preview.Server)
}
if preview.Username != "alice" {
t.Fatalf("username = %q", preview.Username)
}
if !preview.HasPassword || !preview.HasToken {
t.Fatalf("expected secret flags, got %#v", preview)
}
}
func TestUpdateStrmAccountMergesConfigWithoutClearingSecrets(t *testing.T) {
svc := testStrmService(t)
ctx := context.Background()
acct, err := svc.CreateStrmAccount(ctx, "openlist", model.StrmProviderOpenList, map[string]string{
"server": "https://list.example.com",
"username": "alice",
"password": "secret",
"token": "tok",
})
if err != nil {
t.Fatalf("create account: %v", err)
}
updated, err := svc.UpdateStrmAccount(ctx, acct.ID, "openlist-renamed", nil, map[string]string{
"server": "https://list.example.com",
"username": "alice",
})
if err != nil {
t.Fatalf("update account: %v", err)
}
if updated.Name != "openlist-renamed" {
t.Fatalf("name = %q", updated.Name)
}
cfg, err := svc.strmAccountConfig(updated)
if err != nil {
t.Fatalf("decode config: %v", err)
}
if cfg["password"] != "secret" {
t.Fatalf("password was cleared: %#v", cfg)
}
if cfg["token"] != "tok" {
t.Fatalf("token was cleared: %#v", cfg)
}
}
func mustJSON(v any) string {
data, err := json.Marshal(v)
if err != nil {
panic(err)
}
return string(data)
}
+21 -15
View File
@@ -29,9 +29,8 @@ const (
// downloadWorker 下载队列 worker:认领 → 解析直链 → 下载 → 落盘。
//
// 采用「批量认领 + 全局并发限流」:一次认领数个任务,用 StrmService 上的全局信号量
// 限制整个进程「同时换直链+下载」的并发数(与 115 换链风控匹配,见 strmDownloadSemCap),
// 同时让下载充分并行。换链走全局令牌桶(QPS=3)兜底,下载走 CDN 不限速。
// 采用「批量认领 + 按网盘类型限流」:一次认领数个任务,115 与 WebDAV/OpenList
// 使用独立信号量(115 固定 3 防风控;OpenList/CloudDrive2 并发由 download_threads 控制)。
func (s *StrmService) downloadWorker(ctx context.Context) {
const claimBatch = 12 // 每次批量认领的任务数
for {
@@ -42,12 +41,6 @@ func (s *StrmService) downloadWorker(ctx context.Context) {
return
default:
}
// 115 风控/限流熔断:冷却期间整体暂停,不给 WAF 续封机会
if left := s.wafCooldownLeft(); left > 0 {
s.log.Debug("下载队列冷却中", zap.Duration("remaining", left))
sleepContext(ctx, left)
continue
}
tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, claimBatch)
if err != nil {
s.log.Warn("claim strm download task failed", zap.Error(err))
@@ -58,24 +51,37 @@ func (s *StrmService) downloadWorker(ctx context.Context) {
sleepContext(ctx, 2*time.Second)
continue
}
// 并发处理本批任务:每个任务先获取全局下载槽位,槽位内部执行换链+下载。
// 信号量与令牌桶双重限速,确保任意时刻并发换链请求不超过安全阈值。
var wg sync.WaitGroup
for i := range tasks {
wg.Add(1)
go func(i int) {
defer wg.Done()
if !s.acquireDownloadSlot(ctx) {
task := &tasks[i]
if task.Provider == model.StrmProvider115 && s.wafCooldownLeft() > 0 {
s.requeueDownloadTask(task)
return
}
defer s.releaseDownloadSlot()
s.processDownloadTask(ctx, &tasks[i])
if !s.acquireDownloadSlot(ctx, task.Provider) {
s.requeueDownloadTask(task)
return
}
defer s.releaseDownloadSlot(task.Provider)
s.processDownloadTask(ctx, task)
}(i)
}
wg.Wait()
}
}
// requeueDownloadTask 把已认领但未实际执行的任务退回 pending,避免长期停留在 running。
func (s *StrmService) requeueDownloadTask(task *model.StrmDownloadTask) {
task.Status = model.StrmTaskPending
task.StartedAt = nil
if err := s.repo.StrmDownload.Update(context.Background(), task); err != nil {
s.log.Warn("requeue strm download task failed", zap.Error(err), zap.String("id", task.ID))
}
}
func (s *StrmService) processDownloadTask(ctx context.Context, task *model.StrmDownloadTask) {
cleanPath := sanitizeLocalPath(task.LocalPath)
if cleanPath != "" && cleanPath != task.LocalPath {
@@ -648,7 +654,7 @@ func sleepContext(ctx context.Context, d time.Duration) {
// ─── 115 风控/限流熔断 ────────────────────────────────────────────────────────
// triggerWAFCooldown 检测到 115 风控/限流后触发全局冷却,冷却期间下载 worker 暂停。
// triggerWAFCooldown 检测到 115 风控/限流后触发冷却,仅影响 115 下载任务。
// 冷却时间取最大值,避免连续触发时缩短等待。
func (s *StrmService) triggerWAFCooldown() {
s.mu.Lock()
+60
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"testing"
"time"
"github.com/ShukeBta/MMTL/internal/model"
"github.com/ShukeBta/MMTL/internal/repository"
@@ -104,3 +105,62 @@ func TestStrmUploadTasksClearAndRetry(t *testing.T) {
t.Fatalf("expected 2 remaining tasks, got %d", count)
}
}
func TestDownloadSemSeparateByProvider(t *testing.T) {
svc := NewStrmService(nil, zap.NewNop(), nil, nil)
svc.initDownloadSems(6)
if cap(svc.downloadSem115) != strm115DownloadSemCap {
t.Fatalf("115 sem cap = %d, want %d", cap(svc.downloadSem115), strm115DownloadSemCap)
}
if cap(svc.downloadSemDAV) != 6 {
t.Fatalf("dav sem cap = %d, want 6", cap(svc.downloadSemDAV))
}
ctx := context.Background()
for i := 0; i < strm115DownloadSemCap; i++ {
if !svc.acquireDownloadSlot(ctx, model.StrmProvider115) {
t.Fatalf("failed to acquire 115 slot %d", i)
}
}
// 115 槽位占满后,OpenList 仍应能独立获取槽位。
if !svc.acquireDownloadSlot(ctx, model.StrmProviderOpenList) {
t.Fatal("openlist slot should remain available while 115 slots are full")
}
for i := 0; i < 5; i++ {
if !svc.acquireDownloadSlot(ctx, model.StrmProviderOpenList) {
t.Fatalf("failed to acquire extra openlist slot %d", i)
}
}
timeoutCtx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
if svc.acquireDownloadSlot(timeoutCtx, model.StrmProviderOpenList) {
t.Fatal("expected openlist slots to be exhausted after 6 acquisitions")
}
}
func TestRequeueDownloadTask(t *testing.T) {
db := newServiceTestDB(t, &model.StrmDownloadTask{})
repos := repository.New(db)
svc := NewStrmService(nil, zap.NewNop(), repos, nil)
now := time.Now()
task := &model.StrmDownloadTask{
Base: model.Base{ID: "dl-running"},
Status: model.StrmTaskRunning,
Provider: model.StrmProviderOpenList,
StartedAt: &now,
}
if err := db.Create(task).Error; err != nil {
t.Fatalf("insert task: %v", err)
}
svc.requeueDownloadTask(task)
var got model.StrmDownloadTask
if err := db.First(&got, "id = ?", task.ID).Error; err != nil {
t.Fatalf("load task: %v", err)
}
if got.Status != model.StrmTaskPending || got.StartedAt != nil {
t.Fatalf("task not requeued: %+v", got)
}
}
+86 -37
View File
@@ -68,7 +68,7 @@ var StrmSettingDefs = map[string]struct {
StrmSettingUploadMeta: {Default: "false", Label: "上传元数据", Kind: "bool", Help: "同步时把本地元数据上传到远端(需网盘支持写入)"},
StrmSettingDeleteDir: {Default: "false", Label: "清理空目录", Kind: "bool", Help: "清理远端已删除的多余 .strm/元数据后,删除空目录"},
Strm115RelayKeySetting: {Default: "", Label: "115 中继授权共享密钥", Kind: "text", Help: "QMediaSync/MQFamily 中继授权的共享 AES 密钥(OAUTH_RELAY_ENCRYPTION_KEY);不配置则中继授权不可用"},
StrmSettingDownloadThreads: {Default: "3", Label: "下载队列线程数", Kind: "number", Help: "元数据下载并发数"},
StrmSettingDownloadThreads: {Default: "6", Label: "下载队列线程数", Kind: "number", Help: "OpenList/CloudDrive2 元数据下载并发数(115 独立限速为 3)"},
StrmSettingUploadThreads: {Default: "2", Label: "上传队列线程数", Kind: "number", Help: "元数据上传并发数"},
}
@@ -91,33 +91,51 @@ type StrmService struct {
oauthSessions map[string]*strm115AuthSession
wafUntil time.Time // 115 风控/限流熔断截止时间(由 mu 保护)
downloadSem chan struct{} // 全局下载并发信号量:限制整个进程同时进行「换直链+下载」的并发数
downloadSem115 chan struct{} // 115 换直链+下载并发上限(风控兜底)
downloadSemDAV chan struct{} // WebDAV/OpenList/CloudDrive2 元数据下载并发上限
downloadSemOnce sync.Once
}
// strmWAFCooldown 检测到 115 风控/限流后下载队列的全局冷却时长。
// strmWAFCooldown 检测到 115 风控/限流后 115 下载任务的全局冷却时长。
const strmWAFCooldown = 3 * time.Minute
// strmDownloadSemCap 全局同时进行「换直链+下载」的并发上限。
// strm115DownloadSemCap 115 同时进行「换直链+下载」的并发上限。
//
// 115 对换直链接口(/open/ufile/downurl)风控极严:过去把全局 QPS 提到 8 或让多
// worker 高并发换链,会瞬时撞上 WAF 返回 405 阻断页并触发 180 秒冷却,反而更慢。
// 因此用信号量把整个进程同时换直链的并发数压到 3,与令牌桶限速共同兜底:
// 宁可下载稍慢,也绝不触发风控。下载本身走 CDN 不限速。
const strmDownloadSemCap = 3
// 因此把 115 换直链并发压到 3,与令牌桶限速共同兜底;OpenList/CloudDrive2 元数据
// 走 WebDAV,不受此限制,并发由 strm.download_threads 控制。
const strm115DownloadSemCap = 3
// ensureDownloadSem 惰性初始化全局共享的下载并发信号量。
func (s *StrmService) ensureDownloadSem() {
// initDownloadSems 按设置初始化 115 与 WebDAV 两套下载并发信号量。
func (s *StrmService) initDownloadSems(davCap int) {
s.downloadSemOnce.Do(func() {
s.downloadSem = make(chan struct{}, strmDownloadSemCap)
if davCap < 1 {
davCap = 1
}
if davCap > 16 {
davCap = 16
}
s.downloadSem115 = make(chan struct{}, strm115DownloadSemCap)
s.downloadSemDAV = make(chan struct{}, davCap)
})
}
func (s *StrmService) downloadSemForProvider(provider string) chan struct{} {
if provider == model.StrmProvider115 {
return s.downloadSem115
}
return s.downloadSemDAV
}
// acquireDownloadSlot 获取一个下载并发槽位(等待/取消安全)。
func (s *StrmService) acquireDownloadSlot(ctx context.Context) bool {
s.ensureDownloadSem()
func (s *StrmService) acquireDownloadSlot(ctx context.Context, provider string) bool {
sem := s.downloadSemForProvider(provider)
if sem == nil {
return false
}
select {
case s.downloadSem <- struct{}{}:
case sem <- struct{}{}:
return true
case <-ctx.Done():
return false
@@ -125,11 +143,12 @@ func (s *StrmService) acquireDownloadSlot(ctx context.Context) bool {
}
// releaseDownloadSlot 释放一个下载并发槽位。
func (s *StrmService) releaseDownloadSlot() {
if s.downloadSem == nil {
func (s *StrmService) releaseDownloadSlot(provider string) {
sem := s.downloadSemForProvider(provider)
if sem == nil {
return
}
<-s.downloadSem
<-sem
}
// NewStrmService constructs the STRM service.
@@ -151,13 +170,14 @@ func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Conta
func (s *StrmService) Start(ctx context.Context) {
s.sync115RelayKey(ctx)
s.recoverInterruptedSyncs(ctx)
downloadThreads := s.strmIntSetting(ctx, StrmSettingDownloadThreads, 3)
downloadThreads := s.strmIntSetting(ctx, StrmSettingDownloadThreads, 6)
if downloadThreads < 1 {
downloadThreads = 1
}
if downloadThreads > 8 {
downloadThreads = 8
if downloadThreads > 16 {
downloadThreads = 16
}
s.initDownloadSems(downloadThreads)
uploadThreads := s.strmIntSetting(ctx, StrmSettingUploadThreads, 2)
if uploadThreads < 1 {
uploadThreads = 1
@@ -237,9 +257,9 @@ func (s *StrmService) strmAccountConfig(acct *model.StrmAccount) (map[string]str
return cfg, nil
}
// mergeEmbyRemoteConfig 对远程 Emby 账号配置做合并式更新:config 中出现的键
// 覆盖写入(敏感键按明文加密),未出现的键保留原密文;显式空字符串=清除。
func (s *StrmService) mergeEmbyRemoteConfig(existing string, config map[string]string) (string, error) {
// mergeStrmAccountConfig 对账号配置做合并式更新:config 中出现的键覆盖写入(敏感键按明文
// 加密),未出现的键保留原密文;显式空字符串=清除。
func (s *StrmService) mergeStrmAccountConfig(existing string, config map[string]string) (string, error) {
out := map[string]string{}
if strings.TrimSpace(existing) != "" {
if err := json.Unmarshal([]byte(existing), &out); err != nil {
@@ -264,6 +284,44 @@ func (s *StrmService) mergeEmbyRemoteConfig(existing string, config map[string]s
return string(data), nil
}
// StrmAccountConfigPreview 返回可安全回显给前端的非敏感配置字段。
type StrmAccountConfigPreview struct {
URL string `json:"url,omitempty"`
Server string `json:"server,omitempty"`
Username string `json:"username,omitempty"`
HasPassword bool `json:"has_password,omitempty"`
HasToken bool `json:"has_token,omitempty"`
HasAPIKey bool `json:"has_api_key,omitempty"`
}
// StrmAccountConfigPreviewOf 从账号配置提取可回显字段(不含密码/令牌明文)。
func (s *StrmService) StrmAccountConfigPreviewOf(acct *model.StrmAccount) StrmAccountConfigPreview {
var out StrmAccountConfigPreview
if acct == nil || strings.TrimSpace(acct.Config) == "" {
return out
}
raw := map[string]string{}
if err := json.Unmarshal([]byte(acct.Config), &raw); err != nil {
return out
}
out.URL = strings.TrimSpace(raw["url"])
out.Server = strings.TrimSpace(raw["server"])
out.Username = strings.TrimSpace(raw["username"])
for _, k := range []string{"password", "token", "api_key"} {
if v, ok := raw[k]; ok && strings.TrimSpace(v) != "" {
switch k {
case "password":
out.HasPassword = true
case "token":
out.HasToken = true
case "api_key":
out.HasAPIKey = true
}
}
}
return out
}
// HasStrmAccountCredential 报告账号是否已配置核心凭据(用于前端展示)。
func HasStrmAccountCredential(acct *model.StrmAccount) bool {
switch acct.Provider {
@@ -320,22 +378,13 @@ func (s *StrmService) UpdateStrmAccount(ctx context.Context, id, name string, en
if enabled != nil {
acct.Enabled = *enabled
}
if len(config) > 0 {
var enc string
var err error
if acct.Provider == model.StrmProviderEmbyRemote {
// 远程 Emby 账号:合并式更新。config 中出现的键覆盖(敏感键按明文
// 加密写入),未出现的键保留原密文——避免编辑「代理开关」时把已
// 保存的地址与凭据清空。
enc, err = s.mergeEmbyRemoteConfig(acct.Config, config)
} else {
enc, err = s.strmAccountConfigJSON(config, true)
if len(config) > 0 {
enc, err := s.mergeStrmAccountConfig(acct.Config, config)
if err != nil {
return nil, err
}
acct.Config = enc
}
if err != nil {
return nil, err
}
acct.Config = enc
}
if err := s.repo.StrmAccount.Update(ctx, acct); err != nil {
return nil, err
}