diff --git a/internal/service/cloud115/client.go b/internal/service/cloud115/client.go index 27effed..277b686 100644 --- a/internal/service/cloud115/client.go +++ b/internal/service/cloud115/client.go @@ -281,7 +281,7 @@ func IsThrottleCode(code int) bool { func isTokenCode(code int) bool { switch code { - case AccessTokenAuthFail, AccessAuthInvalid, AccessTokenExpiryCode, RefreshTokenInvalid: + case AccessTokenAuthFail, AccessAuthInvalid, AccessTokenExpiryCode, AccessTokenFormatInvalid, RefreshTokenInvalid: return true } return false diff --git a/internal/service/cloud115/const.go b/internal/service/cloud115/const.go index dd92376..c42231f 100644 --- a/internal/service/cloud115/const.go +++ b/internal/service/cloud115/const.go @@ -21,13 +21,14 @@ var ( const ( // 业务错误码 - AccessTokenAuthFail = 40140126 // 访问过期,需刷新 - AccessTokenExpiryCode = 40140125 // 访问过期,需刷新 - AccessAuthInvalid = 40140124 // 访问无效,需刷新 - RefreshTokenInvalid = 40140116 // 需重新授权 - TokenRefreshFail = 40140121 // 刷新失败,可重试 - RequestMaxLimitCode = 770004 // 访问频率过高 - RequestRateLimitCode = 406 // 达到访问上限 + AccessTokenAuthFail = 40140126 // 访问过期,需刷新 + AccessTokenExpiryCode = 40140125 // 访问过期,需刷新 + AccessAuthInvalid = 40140124 // 访问无效,需刷新 + AccessTokenFormatInvalid = 40140123 // access_token 格式错误,需刷新 + RefreshTokenInvalid = 40140116 // 需重新授权 + TokenRefreshFail = 40140121 // 刷新失败,可重试 + RequestMaxLimitCode = 770004 // 访问频率过高 + RequestRateLimitCode = 406 // 达到访问上限 // 刷新 token 的错误码 RefreshTokenFormatInvalid = 40140114 diff --git a/internal/service/strm_sync.go b/internal/service/strm_sync.go index 7197618..e3af7ef 100644 --- a/internal/service/strm_sync.go +++ b/internal/service/strm_sync.go @@ -483,7 +483,8 @@ func cleanDirRel(rel string) string { // 极大地降低 API 请求次数并支持毫秒级/秒级增量同步。 func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error { defer st.flushPendingDownloads() - ctx := st.ctx + ctx, cancel := context.WithCancel(st.ctx) + defer cancel() rootCID := strings.TrimSpace(st.p.RemotePath) if rootCID == "" { rootCID = "0" @@ -618,6 +619,8 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error { dirWorkers = 8 doneDirs atomic.Int64 totalDirs = len(pidList) + errMu sync.Mutex + firstErr error ) if len(pidList) < dirWorkers { dirWorkers = len(pidList) @@ -641,10 +644,19 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error { } detail, err := open115.GetFsDetailByCid(ctx, pid) if err != nil { - st.s.log.Warn("115: 获取目录详情失败", zap.String("pid", pid), zap.Error(err)) // 目录详情解析失败会导致下游文件 rel 无法还原真实父路径, - // seen key 与磁盘路径对不上,增量 prune 会误删本地文件,标记本次扫描不完整。 + // seen key 与磁盘路径对不上:增量 prune 会误删本地文件、上传会 + // 误传本地未变文件、下载会重复下载。这里不是降级容错,而是 + // 直接中止整个同步——宁可本次同步失败,也不带着损坏的相对路径 + // 继续执行造成大规模误删/误传/重下(参考用户反馈"云盘没动却重下重传")。 + errMu.Lock() + if firstErr == nil { + firstErr = fmt.Errorf("115: 解析目录树失败(file_id=%s):%w", pid, err) + } + errMu.Unlock() st.scanIncomplete.Store(true) + cancel() + return } else if detail != nil { // 解析相对路径 relPath := cleanDirRel(detail.RelativePath(rootCID)) @@ -681,6 +693,12 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error { }() } pwg.Wait() + if firstErr != nil { + // 目录树解析失败会导致 rel 塌缩,若继续处理会让大量本地文件 + // 被错误判定为"云端不存在"而重复下载/上传,并可能误删本地文件。 + // 中止本次同步,避免在损坏的相对路径上执行任何写操作。 + return firstErr + } } st.updateSyncMessage(fmt.Sprintf("正在生成 STRM 与同步文件 (共 %d 个)...", len(allFiles))) @@ -698,12 +716,10 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error { if parentVal, ok := st.dirCache.Load(f.Pid); ok && parentVal.(string) != "" { rel = cleanDirRel(parentVal.(string)) + "/" + cleanName } else { - // 父目录不在目录缓存(目录详情先前解析失败),无法还原真实相对路径。 - // 该文件会落到根/错误路径,seen key 与磁盘路径不符,增量 prune 会误删,标记扫描不完整。 - if st.syncType == model.StrmSyncTypeIncremental { - st.scanIncomplete.Store(true) - } - rel = cleanName + // 父目录不在目录缓存,无法还原真实相对路径。若继续用塌缩后的 + // 根路径处理,该文件会被错误判定,导致重复下载/上传或误删本地文件。 + // 目录树不完整时宁可中止本次同步,也不带着损坏的 rel 继续执行。 + return fmt.Errorf("115: 文件 %s 的父目录未解析成功,目录树不完整,中止同步以防误删/误传", cleanName) } } entry := cloud.FileEntry{ diff --git a/internal/service/strm_sync_test.go b/internal/service/strm_sync_test.go index 8d87a8d..c1da6b8 100644 --- a/internal/service/strm_sync_test.go +++ b/internal/service/strm_sync_test.go @@ -2,6 +2,8 @@ package service import ( "context" + "net/http" + "net/http/httptest" "net/url" "os" "path/filepath" @@ -19,6 +21,7 @@ import ( "github.com/ShukeBta/MMTL/internal/model" "github.com/ShukeBta/MMTL/internal/repository" "github.com/ShukeBta/MMTL/internal/service/cloud" + "github.com/ShukeBta/MMTL/internal/service/cloud115" ) // testStrmService 构建带内存库的 StrmService。 @@ -35,10 +38,10 @@ func testStrmService(t *testing.T) *StrmService { sqlDB.SetMaxOpenConns(4) t.Cleanup(func() { _ = sqlDB.Close() }) } - if err := db.AutoMigrate(&model.StrmAccount{}, &model.StrmSyncPath{}, &model.StrmSyncRecord{}, - &model.StrmDownloadTask{}, &model.StrmUploadTask{}, &model.StrmDirCache{}, &model.Setting{}); err != nil { - t.Fatal(err) - } + if err := db.AutoMigrate(&model.StrmAccount{}, &model.StrmSyncPath{}, &model.StrmSyncRecord{}, + &model.StrmDownloadTask{}, &model.StrmUploadTask{}, &model.StrmDirCache{}, &model.Setting{}); err != nil { + t.Fatal(err) + } repos := repository.New(db) ctx := context.Background() if err := repos.Setting.Set(ctx, StrmSettingBaseURL, "http://test.local:8096"); err != nil { @@ -210,7 +213,6 @@ func TestStrmFullAndIncrementalSync(t *testing.T) { } } - // TestStrmCronMatches cron 表达式匹配。 func TestStrmCronMatches(t *testing.T) { cases := []struct { @@ -494,142 +496,215 @@ func TestWalkRemoteConcurrent(t *testing.T) { if walkErr != nil { t.Fatal(walkErr) } - if strmCount != 5 { - t.Errorf("生成的 .strm 数量 = %d,期望 5", strmCount) - } + if strmCount != 5 { + t.Errorf("生成的 .strm 数量 = %d,期望 5", strmCount) + } +} + +// TestStrmBatchEnqueueAndConcurrentClaim 测试大规模批量入库及多协程并发认领无死锁 +func TestStrmBatchEnqueueAndConcurrentClaim(t *testing.T) { + svc := testStrmService(t) + ctx := context.Background() + + // 1. 批量插入 200 个下载任务 + tasks := make([]*model.StrmDownloadTask, 0, 200) + for i := 0; i < 200; i++ { + tasks = append(tasks, &model.StrmDownloadTask{ + SyncPathID: "test-sync-path", + AccountID: "test-acct", + Provider: model.StrmProvider115, + FileName: filepath.Base(string(rune('a'+i%26))) + ".nfo", + LocalPath: filepath.Join(t.TempDir(), string(rune('a'+i%26)), "test.nfo"), + Status: model.StrmTaskPending, + }) + } + if err := svc.repo.StrmDownload.CreateInBatches(ctx, tasks, 50); err != nil { + t.Fatalf("CreateInBatches failed: %v", err) } - // TestStrmBatchEnqueueAndConcurrentClaim 测试大规模批量入库及多协程并发认领无死锁 - func TestStrmBatchEnqueueAndConcurrentClaim(t *testing.T) { - svc := testStrmService(t) - ctx := context.Background() + // 2. 验证 ActiveLocalPathMap + activeMap, err := svc.repo.StrmDownload.GetActiveLocalPathMap(ctx, "test-sync-path") + if err != nil { + t.Fatalf("GetActiveLocalPathMap failed: %v", err) + } + if len(activeMap) == 0 { + t.Fatal("expected active local path map to have entries") + } - // 1. 批量插入 200 个下载任务 - tasks := make([]*model.StrmDownloadTask, 0, 200) - for i := 0; i < 200; i++ { - tasks = append(tasks, &model.StrmDownloadTask{ - SyncPathID: "test-sync-path", - AccountID: "test-acct", - Provider: model.StrmProvider115, - FileName: filepath.Base(string(rune('a'+i%26))) + ".nfo", - LocalPath: filepath.Join(t.TempDir(), string(rune('a'+i%26)), "test.nfo"), - Status: model.StrmTaskPending, - }) - } - if err := svc.repo.StrmDownload.CreateInBatches(ctx, tasks, 50); err != nil { - t.Fatalf("CreateInBatches failed: %v", err) - } - - // 2. 验证 ActiveLocalPathMap - activeMap, err := svc.repo.StrmDownload.GetActiveLocalPathMap(ctx, "test-sync-path") - if err != nil { - t.Fatalf("GetActiveLocalPathMap failed: %v", err) - } - if len(activeMap) == 0 { - t.Fatal("expected active local path map to have entries") - } - - // 3. 模拟 6 个 worker 并发 ClaimPendingDownload - claimedCount := 0 - var claimMu sync.Mutex - var wg sync.WaitGroup - for w := 0; w < 6; w++ { - wg.Add(1) - go func() { - defer wg.Done() - for { - batch, err := svc.repo.StrmDownload.ClaimPendingDownload(ctx, 10) - if err != nil { - t.Errorf("concurrent ClaimPendingDownload failed: %v", err) - return - } - if len(batch) == 0 { - return - } - claimMu.Lock() - claimedCount += len(batch) - claimMu.Unlock() + // 3. 模拟 6 个 worker 并发 ClaimPendingDownload + claimedCount := 0 + var claimMu sync.Mutex + var wg sync.WaitGroup + for w := 0; w < 6; w++ { + wg.Add(1) + go func() { + defer wg.Done() + for { + batch, err := svc.repo.StrmDownload.ClaimPendingDownload(ctx, 10) + if err != nil { + t.Errorf("concurrent ClaimPendingDownload failed: %v", err) + return } - }() - } - wg.Wait() - - if claimedCount != 200 { - t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount) + if len(batch) == 0 { + return + } + claimMu.Lock() + claimedCount += len(batch) + claimMu.Unlock() } - } + }() + } + wg.Wait() - // TestStrmDuplicateFileConflictResolution 测试远端存在多个同名不同大小文件时,本地确定性仲裁,避免增量死循环 - func TestStrmDuplicateFileConflictResolution(t *testing.T) { - svc := testStrmService(t) - localDir := t.TempDir() + if claimedCount != 200 { + t.Fatalf("expected all 200 tasks claimed, got %d", claimedCount) + } +} - p := &model.StrmSyncPath{ - Base: model.Base{ID: "dup-test-path"}, - Provider: model.StrmProvider115, - RemotePath: "root", - LocalPath: localDir, - DownloadMeta: true, - } +// TestStrmDuplicateFileConflictResolution 测试远端存在多个同名不同大小文件时,本地确定性仲裁,避免增量死循环 +func TestStrmDuplicateFileConflictResolution(t *testing.T) { + svc := testStrmService(t) + localDir := t.TempDir() - 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) - } + 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) + } + +} + +// TestWalk115FlatAbortsOnDirResolveFailure 回归测试:115 开放平台 token 失效/目录详情 +// 解析失败时,同步必须中止而不是带着塌缩的 rel 继续处理,否则会导致本地大量元数据 +// 被误判为"云端不存在"而重复下载/上传,甚至误删本地文件(用户反馈"云盘没动却重下重传")。 +func TestWalk115FlatAbortsOnDirResolveFailure(t *testing.T) { + svc := testStrmService(t) + localDir := t.TempDir() + + acct := &model.StrmAccount{ + Name: "fake115", + Provider: "cloud115", + Config: "{}", + Enabled: true, + } + if err := svc.repo.StrmAccount.Create(context.Background(), acct); err != nil { + t.Fatal(err) + } + p := &model.StrmSyncPath{ + Base: model.Base{ID: "abort-path"}, + AccountID: acct.ID, + Provider: model.StrmProvider115, + RemotePath: "0", + LocalPath: localDir, + } + + // 115 mock:文件列表返回一个视频(父目录 999 不在缓存,需要 get_info), + // get_info 恒返回 access_token 格式错误(40140123)→ 目录树解析失败。 + api := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/open/ufile/files": + w.Write([]byte(`{"state":true,"count":1,"data":[{"fid":"100","pid":"999","fc":1,"fn":"movie.mkv","pc":"pc1","upt":1700000000,"fs":1024}]}`)) + case "/open/folder/get_info": + w.Write([]byte(`{"state":false,"code":40140123,"message":"access_token 格式错误"}`)) + default: + t.Errorf("unexpected path %s", r.URL.Path) + } + })) + defer api.Close() + + oldPro := cloud115.ProAPIBase + cloud115.ProAPIBase = api.URL + defer func() { cloud115.ProAPIBase = oldPro }() + + oc := cloud115.NewOpenClient("app", "at", "rt") + st := &strmSyncState{ + s: svc, + ctx: context.Background(), + p: p, + provider: cloud.NewOpenAPI115("app", "at", "rt"), + cfg: &strmPathConfig{VideoExt: []string{"mkv"}, MetaExt: []string{"nfo"}, AddPath: 1, DownloadMeta: false}, + rec: &model.StrmSyncRecord{}, + syncType: model.StrmSyncTypeFull, + dirCache: sync.Map{}, + seenVideo: map[string]bool{}, + seenMeta: map[string]bool{}, + remoteMeta: map[string]int64{}, + } + err := st.walk115Flat(oc) + if err == nil { + t.Fatal("expected walk115Flat to abort on dir-resolve failure, got nil error") + } + + // 中止后不允许产生任何部分写入(本地不允许生成 .strm 文件)。 + var strmCount int + _ = filepath.WalkDir(localDir, func(path string, d os.DirEntry, err error) error { + if err == nil && !d.IsDir() && strings.HasSuffix(d.Name(), ".strm") { + strmCount++ + } + return nil + }) + if strmCount != 0 { + t.Fatalf("expected no .strm written after abort, got %d", strmCount) + } +}