Compare commits

...

10 Commits

Author SHA1 Message Date
github-actions[bot] 60c815a8b3 chore: bump version to 0.0.38 [skip ci] 2026-08-26 09:06:04 +00:00
truewhile 41b155ea31 Merge branch 'main' of https://github.com/truewhile/MMTL 2026-08-26 17:05:45 +08:00
truewhile 2888ae8bf7 8 2026-08-26 17:05:41 +08:00
github-actions[bot] 7363064d89 chore: bump version to 0.0.37 [skip ci] 2026-08-26 08:16:56 +00:00
truewhile 1d53bf2ae1 7 2026-08-26 16:16:39 +08:00
github-actions[bot] 618165ec31 chore: bump version to 0.0.36 [skip ci] 2026-08-26 07:51:14 +00:00
truewhile 87c66a9b8c 6 2026-08-26 15:50:59 +08:00
github-actions[bot] 1ea4724261 chore: bump version to 0.0.35 [skip ci] 2026-08-26 06:50:29 +00:00
truewhile 3d372f039e Merge branch 'main' of https://github.com/truewhile/MMTL 2026-08-26 14:50:10 +08:00
truewhile 0384017e98 6 2026-08-26 14:50:06 +08:00
13 changed files with 312 additions and 63 deletions
+1 -1
View File
@@ -1 +1 @@
0.0.34
0.0.38
+38 -8
View File
@@ -16,12 +16,13 @@ import (
)
type createLibraryReq struct {
Name string `json:"name" binding:"required"`
Path string `json:"path"`
Paths []string `json:"paths"`
Roots []service.LibraryRootInput `json:"roots"`
Type string `json:"type"`
CoverURL string `json:"cover_url"`
Name string `json:"name"`
Path string `json:"path"`
Paths []string `json:"paths"`
Roots []service.LibraryRootInput `json:"roots"`
Type string `json:"type"`
CoverURL string `json:"cover_url"`
CreatePerSubfolder bool `json:"create_per_subfolder"`
}
func listLibrariesHandler(svc *service.Container) gin.HandlerFunc {
@@ -88,9 +89,38 @@ func createLibraryHandler(svc *service.Container) gin.HandlerFunc {
}
}
if len(roots) == 0 && strings.TrimSpace(req.Path) != "" {
roots = append(roots, service.LibraryRootInput{Path: req.Path})
roots = append(roots, service.LibraryRootInput{Path: req.Path})
}
var l *model.Library
if req.CreatePerSubfolder {
parent := ""
if len(roots) > 0 {
parent = roots[0].Path
} else if strings.TrimSpace(req.Path) != "" {
parent = req.Path
}
l, err := svc.Media.CreateLibraryWithRootsAndCover(c.Request.Context(), req.Name, req.Type, req.CoverURL, roots)
created, err := svc.Media.CreateLibrariesPerSubfolder(c.Request.Context(), parent, req.Type, req.CoverURL)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
uid, _ := c.Get("ctx_user_id")
for i := range created {
lib := &created[i]
svc.Audit.Record(c.Request.Context(), toString(uid), "library.create", lib.ID, c.ClientIP(), lib.Path)
if svc.Watcher != nil {
go func() { _ = svc.Watcher.Refresh(context.Background()) }()
}
for _, root := range lib.Roots {
if root.Enabled {
queueLibraryRootScan(svc, lib.ID, root.ID)
}
}
}
c.JSON(http.StatusCreated, gin.H{"libraries": created})
return
}
l, err := svc.Media.CreateLibraryWithRootsAndCover(c.Request.Context(), req.Name, req.Type, req.CoverURL, roots)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
+7 -2
View File
@@ -26,10 +26,15 @@ var (
executorOnce sync.Once
)
// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=2, QPM=120, QPH=6000,保障 115 API 调用安全不超频)。
// GetGlobalExecutor 获取全局队列执行器单例(默认 QPS=3, QPM=200, QPH=12000,保障 115 API 调用安全不超频)。
//
// 历史教训:QPS 提到 8 后,下载换直链接口(/open/ufile/downurl,WAF 重点盯防对象)
// 瞬时突发撞上 115 风控,返回阿里云 405 阻断页(HTTP 405),导致全量同步失败。
// 因此回调到 3——这是经过实测的安全上限:宁慢勿触发风控,一旦 405 冷却 180 秒,
// 整体吞吐反而更低。下载实际走 CDN 不受此限速影响,瓶颈仅在换链环节。
func GetGlobalExecutor() *QueueExecutor {
executorOnce.Do(func() {
globalExecutor = NewQueueExecutor(2, 120, 6000)
globalExecutor = NewQueueExecutor(3, 200, 12000)
})
return globalExecutor
}
+41
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
@@ -59,6 +60,46 @@ func (s *MediaService) CreateLibraryWithRootsAndCover(ctx context.Context, name,
return lib, nil
}
// CreateLibrariesPerSubfolder 为 parent 目录下的每个直接子目录各建一个媒体库,
// 媒体库名取子目录名,路径指向该子目录。kind 为空时按子目录名推断类型。
func (s *MediaService) CreateLibrariesPerSubfolder(ctx context.Context, parent, kind, coverURL string) ([]model.Library, error) {
parent = strings.TrimSpace(parent)
if parent == "" {
return nil, errors.New("parent path required")
}
dir, err := resolveAccessibleLibraryPath(parent)
if err != nil {
return nil, err
}
entries, err := os.ReadDir(dir)
if err != nil {
return nil, fmt.Errorf("read directory failed: %w", err)
}
subdirs := make([]string, 0, len(entries))
for _, entry := range entries {
if !entry.IsDir() {
continue
}
if strings.HasPrefix(entry.Name(), ".") {
continue
}
subdirs = append(subdirs, filepath.Join(dir, entry.Name()))
}
if len(subdirs) == 0 {
return nil, errors.New("no subfolders found")
}
created := make([]model.Library, 0, len(subdirs))
for _, subdir := range subdirs {
name := filepath.Base(subdir)
lib, err := s.CreateLibraryWithRootsAndCover(ctx, name, kind, coverURL, []LibraryRootInput{{Path: subdir}})
if err != nil {
return nil, fmt.Errorf("create library for %s: %w", subdir, err)
}
created = append(created, *lib)
}
return created, nil
}
func (s *MediaService) UpdateLibraryCover(ctx context.Context, libraryID, coverURL string) error {
return s.repo.DB.WithContext(ctx).Model(&model.Library{}).Where("id = ?", libraryID).
Update("cover_url", strings.TrimSpace(coverURL)).Error
+20 -2
View File
@@ -13,6 +13,7 @@ import (
"os"
"path/filepath"
"strings"
"sync"
"time"
"go.uber.org/zap"
@@ -27,7 +28,12 @@ const (
)
// downloadWorker 下载队列 worker:认领 → 解析直链 → 下载 → 落盘。
//
// 采用「批量认领 + 全局并发限流」:一次认领数个任务,用 StrmService 上的全局信号量
// 限制整个进程「同时换直链+下载」的并发数(与 115 换链风控匹配,见 strmDownloadSemCap),
// 同时让下载充分并行。换链走全局令牌桶(QPS=3)兜底,下载走 CDN 不限速。
func (s *StrmService) downloadWorker(ctx context.Context) {
const claimBatch = 12 // 每次批量认领的任务数
for {
select {
case <-ctx.Done():
@@ -42,7 +48,7 @@ func (s *StrmService) downloadWorker(ctx context.Context) {
sleepContext(ctx, left)
continue
}
tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, 1)
tasks, err := s.repo.StrmDownload.ClaimPendingDownload(ctx, claimBatch)
if err != nil {
s.log.Warn("claim strm download task failed", zap.Error(err))
sleepContext(ctx, 3*time.Second)
@@ -52,9 +58,21 @@ func (s *StrmService) downloadWorker(ctx context.Context) {
sleepContext(ctx, 2*time.Second)
continue
}
// 并发处理本批任务:每个任务先获取全局下载槽位,槽位内部执行换链+下载。
// 信号量与令牌桶双重限速,确保任意时刻并发换链请求不超过安全阈值。
var wg sync.WaitGroup
for i := range tasks {
s.processDownloadTask(ctx, &tasks[i])
wg.Add(1)
go func(i int) {
defer wg.Done()
if !s.acquireDownloadSlot(ctx) {
return
}
defer s.releaseDownloadSlot()
s.processDownloadTask(ctx, &tasks[i])
}(i)
}
wg.Wait()
}
}
+37
View File
@@ -90,11 +90,48 @@ type StrmService struct {
running map[string]context.CancelFunc // sync path id -> cancel
oauthSessions map[string]*strm115AuthSession
wafUntil time.Time // 115 风控/限流熔断截止时间(由 mu 保护)
downloadSem chan struct{} // 全局下载并发信号量:限制整个进程同时进行「换直链+下载」的并发数
downloadSemOnce sync.Once
}
// strmWAFCooldown 检测到 115 风控/限流后下载队列的全局冷却时长。
const strmWAFCooldown = 3 * time.Minute
// strmDownloadSemCap 全局同时进行「换直链+下载」的并发上限。
//
// 115 对换直链接口(/open/ufile/downurl)风控极严:过去把全局 QPS 提到 8 或让多
// worker 高并发换链,会瞬时撞上 WAF 返回 405 阻断页并触发 180 秒冷却,反而更慢。
// 因此用信号量把整个进程同时换直链的并发数压到 3,与令牌桶限速共同兜底:
// 宁可下载稍慢,也绝不触发风控。下载本身走 CDN 不限速。
const strmDownloadSemCap = 3
// ensureDownloadSem 惰性初始化全局共享的下载并发信号量。
func (s *StrmService) ensureDownloadSem() {
s.downloadSemOnce.Do(func() {
s.downloadSem = make(chan struct{}, strmDownloadSemCap)
})
}
// acquireDownloadSlot 获取一个下载并发槽位(等待/取消安全)。
func (s *StrmService) acquireDownloadSlot(ctx context.Context) bool {
s.ensureDownloadSem()
select {
case s.downloadSem <- struct{}{}:
return true
case <-ctx.Done():
return false
}
}
// releaseDownloadSlot 释放一个下载并发槽位。
func (s *StrmService) releaseDownloadSlot() {
if s.downloadSem == nil {
return
}
<-s.downloadSem
}
// NewStrmService constructs the STRM service.
func NewStrmService(cfg *config.Config, log *zap.Logger, repos *repository.Container, crypto *CryptoService) *StrmService {
return &StrmService{
+41 -11
View File
@@ -40,8 +40,10 @@ type strmSyncState struct {
seenVideo map[string]bool // "v:"+去掉扩展名的相对路径 → 远端存在该视频
seenMeta map[string]bool // "m:"+相对路径 → 远端存在该元数据
remoteMeta map[string]int64 // 远端元数据大小(上传比对用)
activeDownloadPaths map[string]bool // 本地已在排队/进行的下载任务路径(内存去重)
activeUploadPaths map[string]bool // 本地已在排队/进行的上传任务路径(内存去重)
seenMetaTarget map[string]cloud.FileEntry
seenVideoTarget map[string]cloud.FileEntry
activeDownloadPaths map[string]bool // 本地已在排队/进行的下载任务路径(内存去重)
activeUploadPaths map[string]bool // 本地已在排队/进行的上传任务路径(内存去重)
pendingDownloads []*model.StrmDownloadTask
pendingUploads []*model.StrmUploadTask
dirCache sync.Map // dirID (string) -> relativePath (string)
@@ -172,15 +174,17 @@ func (s *StrmService) runSync(ctx context.Context, p *model.StrmSyncPath, rec *m
return
}
st := &strmSyncState{
s: s,
ctx: ctx,
p: p,
cfg: cfg,
rec: rec,
syncType: rec.SyncType,
seenVideo: map[string]bool{},
seenMeta: map[string]bool{},
remoteMeta: map[string]int64{},
s: s,
ctx: ctx,
p: p,
cfg: cfg,
rec: rec,
syncType: rec.SyncType,
seenVideo: map[string]bool{},
seenMeta: map[string]bool{},
remoteMeta: map[string]int64{},
seenMetaTarget: map[string]cloud.FileEntry{},
seenVideoTarget: map[string]cloud.FileEntry{},
}
if p.Provider != model.StrmProviderLocal {
acct, err := s.repo.StrmAccount.FindByID(ctx, p.AccountID)
@@ -686,6 +690,18 @@ func (st *strmSyncState) handleVideo(entry cloud.FileEntry, rel, ext string) {
return
}
st.mu.Lock()
if st.seenVideoTarget == nil {
st.seenVideoTarget = map[string]cloud.FileEntry{}
}
if _, exists := st.seenVideoTarget[target]; exists {
st.mu.Unlock()
st.touchProgress()
return
}
st.seenVideoTarget[target] = entry
st.mu.Unlock()
// 增量同步模式快速检查:本地 strm 文件存在、非空且修改时间与远端 mtime 一致,直接跳过无需读磁盘
if st.syncType == model.StrmSyncTypeIncremental && entry.MTime > 0 {
if info, err := os.Stat(target); err == nil && info.Size() > 0 && info.ModTime().Unix() == entry.MTime {
@@ -835,6 +851,20 @@ func (st *strmSyncState) handleMeta(entry cloud.FileEntry, rel, ext string) {
if err != nil {
return
}
st.mu.Lock()
if st.seenMetaTarget == nil {
st.seenMetaTarget = map[string]cloud.FileEntry{}
}
if _, exists := st.seenMetaTarget[target]; exists {
// 该本地目标路径在当前批次中已被处理(存在同名/重名冲突),直接忽略重复项,避免多份不同大小的文件在本地交替覆盖导致增量死循环
st.mu.Unlock()
st.touchProgress()
return
}
st.seenMetaTarget[target] = entry
st.mu.Unlock()
if info, err := os.Stat(target); err == nil && info.Size() == entry.Size {
st.touchProgress()
return
+76 -2
View File
@@ -554,8 +554,82 @@ func TestWalkRemoteConcurrent(t *testing.T) {
}
wg.Wait()
if claimedCount != 200 {
t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount)
if claimedCount != 200 {
t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount)
}
}
// TestStrmDuplicateFileConflictResolution 测试远端存在多个同名不同大小文件时,本地确定性仲裁,避免增量死循环
func TestStrmDuplicateFileConflictResolution(t *testing.T) {
svc := testStrmService(t)
localDir := t.TempDir()
p := &model.StrmSyncPath{
Base: model.Base{ID: "dup-test-path"},
Provider: model.StrmProvider115,
RemotePath: "root",
LocalPath: localDir,
DownloadMeta: true,
}
st := &strmSyncState{
s: svc,
ctx: context.Background(),
p: p,
cfg: &strmPathConfig{DownloadMeta: true, MetaExt: []string{"nfo"}},
rec: &model.StrmSyncRecord{},
seenMeta: map[string]bool{},
remoteMeta: map[string]int64{},
seenMetaTarget: map[string]cloud.FileEntry{},
seenVideoTarget: map[string]cloud.FileEntry{},
}
// 模拟远端同目录下存在两个同名不同大小的 nfo 文件 (115 历史重复上传)
// entry1: 较早文件 (MTime: 1000, Size: 100)
entry1 := cloud.FileEntry{ID: "f1", Name: "test.nfo", Size: 100, MTime: 1000, PickCode: "p1"}
// entry2: 较新文件 (MTime: 2000, Size: 200)
entry2 := cloud.FileEntry{ID: "f2", Name: "test.nfo", Size: 200, MTime: 2000, PickCode: "p2"}
// 第一次全量处理:两者都在列表中
st.handleMeta(entry1, "test.nfo", ".nfo")
st.handleMeta(entry2, "test.nfo", ".nfo")
st.flushPendingDownloads()
// 验证仲裁结果:最终只产生 1 个下载任务,且使用的是首个匹配项 (Size 100/p1)
tasks, _, err := svc.repo.StrmDownload.List(context.Background(), "", 1, 10)
if err != nil {
t.Fatal(err)
}
if len(tasks) != 1 {
t.Fatalf("expected 1 download task after conflict resolution, got %d", len(tasks))
}
if tasks[0].Size != 100 || tasks[0].RemoteRef != "p1" {
t.Fatalf("expected task with size 100/p1, got size=%d ref=%s", tasks[0].Size, tasks[0].RemoteRef)
}
// 模拟该任务下载落盘完成
writeFile(t, filepath.Join(localDir, "test.nfo"), strings.Repeat("x", 100))
// 第二次增量同步:两者再次依次扫描
st2 := &strmSyncState{
s: svc,
ctx: context.Background(),
p: p,
cfg: &strmPathConfig{DownloadMeta: true, MetaExt: []string{"nfo"}},
rec: &model.StrmSyncRecord{},
seenMeta: map[string]bool{},
remoteMeta: map[string]int64{},
seenMetaTarget: map[string]cloud.FileEntry{},
seenVideoTarget: map[string]cloud.FileEntry{},
}
st2.handleMeta(entry1, "test.nfo", ".nfo")
st2.handleMeta(entry2, "test.nfo", ".nfo")
st2.flushPendingDownloads()
// 验证:不会新增任何下载任务,NewMeta 为 0,增量跳过
if st2.rec.NewMeta != 0 {
t.Fatalf("expected 0 new meta on incremental sync, got %d", st2.rec.NewMeta)
}
}
-25
View File
@@ -1,25 +0,0 @@
package main
import (
"context"
"fmt"
"github.com/ShukeBta/MMTL/internal/service"
"github.com/ShukeBta/MMTL/internal/service/cloud115"
"go.uber.org/zap"
)
func main() {
crypto := service.NewCryptoService("test-secret", zap.NewNop())
// Let's test with a mock RemoteFileDetail
d := &cloud115.RemoteFileDetail{
FileId: "3251154147730910635",
FileName: "出包王女",
Paths: []struct {
FileId string
Name string
}{
{FileId: "0", Name: "根目录"},
{FileId: "3238787832374488117", Name: "影视库"},
{FileId: "3238787913223892116", Name: "动漫"},
},
}
fmt.Println("RelativePath when rootCID is 3238787832374488117:", d.RelativePath("3238787832374488117"))
}
+3
View File
@@ -102,6 +102,9 @@ export const libraryAPI = {
createWithRoots: (name: string, type: string, roots: LibraryRootInput[], coverURL = '') =>
api.post<Library>('/libraries', { name, type, roots, cover_url: coverURL }).then((r) => r.data),
createPerSubfolder: (parentPath: string, type: string, coverURL = '') =>
api.post<{ libraries: Library[] }>('/libraries', { path: parentPath, type, cover_url: coverURL, create_per_subfolder: true }).then((r) => r.data),
update: (id: string, payload: { cover_url: string }) =>
api.patch<Library>(`/libraries/${id}`, payload).then((r) => r.data),
+2
View File
@@ -13,9 +13,11 @@ export function AdminLibraryPanel() {
type={createForm.type}
coverURL={createForm.coverURL}
roots={createForm.roots}
createPerSubfolder={createForm.createPerSubfolder}
onNameChange={createForm.setName}
onTypeChange={createForm.setType}
onCoverURLChange={createForm.setCoverURL}
onCreatePerSubfolderChange={createForm.setCreatePerSubfolder}
onRootChange={createForm.updateRoot}
onAddRoot={createForm.addRoot}
onRemoveRoot={createForm.removeRoot}
+26 -6
View File
@@ -9,9 +9,11 @@ type CreateFormProps = {
type: string
coverURL: string
roots: RootDraft[]
createPerSubfolder: boolean
onNameChange: (value: string) => void
onTypeChange: (value: string) => void
onCoverURLChange: (value: string) => void
onCreatePerSubfolderChange: (value: boolean) => void
onRootChange: (index: number, patch: Partial<RootDraft>) => void
onAddRoot: () => void
onRemoveRoot: (index: number) => void
@@ -23,9 +25,11 @@ export function AdminLibraryCreateForm({
type,
coverURL,
roots,
createPerSubfolder,
onNameChange,
onTypeChange,
onCoverURLChange,
onCreatePerSubfolderChange,
onRootChange,
onAddRoot,
onRemoveRoot,
@@ -50,9 +54,9 @@ export function AdminLibraryCreateForm({
<>
<form onSubmit={onSubmit} className="glass-panel grid gap-3 md:grid-cols-4">
<input
required
required={!createPerSubfolder}
className="input-base"
placeholder="名称"
placeholder={createPerSubfolder ? '父级媒体库名(批量模式忽略)' : '名称'}
value={name}
onChange={(e) => onNameChange(e.target.value)}
/>
@@ -81,15 +85,31 @@ export function AdminLibraryCreateForm({
onRemove={onRemoveRoot}
/>
))}
<button type="button" className="inline-flex items-center gap-2 rounded-lg border px-3 py-2 text-sm" onClick={onAddRoot}>
<Plus size={16} /> 添加路径
</button>
{!createPerSubfolder && (
<button type="button" className="inline-flex items-center gap-2 rounded-lg border px-3 py-2 text-sm" onClick={onAddRoot}>
<Plus size={16} /> 添加路径
</button>
)}
</div>
<p className="md:col-span-4 -mt-2 text-xs text-sand-500">
支持直接点选或手动输入;名称和类型与现有媒体库一致时,会自动把这里填写的路径追加到该媒体库。
</p>
<label className="md:col-span-4 flex items-center gap-2 text-sm text-ink-100">
<input
type="checkbox"
className="h-4 w-4 accent-brand-400"
checked={createPerSubfolder}
onChange={(e) => onCreatePerSubfolderChange(e.target.checked)}
/>
<span>按目录下每个子文件夹各建一个媒体库(媒体库名取子文件夹名)</span>
</label>
{createPerSubfolder && (
<p className="md:col-span-4 -mt-2 text-xs text-sand-500">
批处理模式:仅取上方第一个路径作为父级目录,会为其中每个子文件夹分别创建媒体库,可自选类型用于整体推断。
</p>
)}
<button type="submit" className="neon-button md:col-span-4">
新建 / 追加路径
{createPerSubfolder ? '按目录批量创建' : '新建 / 追加路径'}
</button>
</form>
+20 -6
View File
@@ -32,20 +32,32 @@ function useCreateLibraryForm(refresh: () => Promise<void>) {
const [roots, setRoots] = useState<RootDraft[]>([emptyRootDraft()])
const [type, setType] = useState('movie')
const [coverURL, setCoverURL] = useState('')
const [createPerSubfolder, setCreatePerSubfolder] = useState(false)
const handleCreate = async (e: FormEvent) => {
e.preventDefault()
try {
const payload = createRootPayload(roots)
if (payload.length === 0) {
toast.error('请至少填写一个路径')
return
if (createPerSubfolder) {
const parentPath = roots[0]?.path?.trim()
if (!parentPath) {
toast.error('请先选择或填写父级目录')
return
}
const { libraries } = await libraryAPI.createPerSubfolder(parentPath, type, coverURL.trim())
toast.success(`已按目录创建 ${libraries.length} 个媒体库`)
} else {
const payload = createRootPayload(roots)
if (payload.length === 0) {
toast.error('请至少填写一个路径')
return
}
await libraryAPI.createWithRoots(name, type, payload, coverURL.trim())
toast.success('媒体库已保存')
}
await libraryAPI.createWithRoots(name, type, payload, coverURL.trim())
toast.success('媒体库已保存')
setName('')
setRoots([emptyRootDraft()])
setCoverURL('')
setCreatePerSubfolder(false)
await refresh()
} catch (err: unknown) {
toast.error(apiErrorMessage(err, '创建失败'))
@@ -61,9 +73,11 @@ function useCreateLibraryForm(refresh: () => Promise<void>) {
type,
coverURL,
roots,
createPerSubfolder,
setName,
setType,
setCoverURL,
setCreatePerSubfolder,
updateRoot,
addRoot: () => setRoots((prev) => [...prev, emptyRootDraft()]),
removeRoot: (index: number) => setRoots((prev) => (prev.length <= 1 ? prev : prev.filter((_, i) => i !== index))),