From 2aeedcc1822931a392c59cf682952ae30bb08849 Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Sun, 6 Sep 2026 00:01:40 +0800 Subject: [PATCH] =?UTF-8?q?bug=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/service/strm_queue.go | 5 +- internal/service/strm_sync.go | 128 +++++++++++------ internal/service/strm_sync_test.go | 213 ++++++++++++++++++++++------- 3 files changed, 255 insertions(+), 91 deletions(-) diff --git a/internal/service/strm_queue.go b/internal/service/strm_queue.go index 4a98486..184ff40 100644 --- a/internal/service/strm_queue.go +++ b/internal/service/strm_queue.go @@ -366,7 +366,7 @@ func (s *StrmService) processUpload115(ctx context.Context, task *model.StrmUplo finish(model.StrmTaskFailed, "该网盘不支持元数据上传") return } - // 以本地为准:网盘端已有同名但内容不同的旧元数据时,先删除旧文件再上传。 + // 以本地为准:网盘端已有同名但内容不同的旧元数据时,先批量删除所有旧副本再上传。 // 115 的上传接口不保证同名覆盖,直接上传可能产生同名重复文件;删除失败则 // 任务重试(旧文件 ID 失效的场景会在下次同步后自动修复)。 if task.RemoteRef != "" { @@ -375,7 +375,8 @@ func (s *StrmService) processUpload115(ctx context.Context, task *model.StrmUplo finish(model.StrmTaskFailed, "该网盘不支持删除远端旧元数据") return } - if err := open115.OpenClient().DeleteFiles(ctx, task.RemotePath, task.RemoteRef); err != nil { + refs := strings.Split(task.RemoteRef, ",") + if err := open115.OpenClient().DeleteFiles(ctx, task.RemotePath, refs...); err != nil { s.uploadTaskFailWithRetry(task, "删除网盘旧元数据失败:"+err.Error()) return } diff --git a/internal/service/strm_sync.go b/internal/service/strm_sync.go index 9790582..15e043e 100644 --- a/internal/service/strm_sync.go +++ b/internal/service/strm_sync.go @@ -25,6 +25,14 @@ import ( "github.com/truewhile/MeBox/internal/service/cloud115" ) +// remoteMetaItem 记录远端存在的单个元数据文件副本信息(大小、文件ID、内容SHA1、修改时间)。 +type remoteMetaItem struct { + ID string + Size int64 + Sha1 string + MTime int64 +} + // strmSyncState 是一次同步执行的上下文。 type strmSyncState struct { s *StrmService @@ -37,13 +45,11 @@ type strmSyncState struct { syncType string mu sync.Mutex - processed int // 已处理文件计数(用于定期落库进度) - lastProgressFlush time.Time // 上次进度落库时间 - seenVideo map[string]bool // "v:"+去掉扩展名的相对路径 → 远端存在该视频 - seenMeta map[string]bool // "m:"+相对路径 → 远端存在该元数据 - remoteMeta map[string]int64 // 远端元数据大小(上传比对用) - remoteMetaRef map[string]string // "m:"+相对路径 → 远端元数据文件引用(115 文件 ID,覆盖上传前删除旧文件用) - remoteMetaSha1 map[string]string // "m:"+相对路径 → 远端元数据内容 SHA1(115 列表返回;上传/下载精确比对用,其他网盘为空) + processed int // 已处理文件计数(用于定期落库进度) + lastProgressFlush time.Time // 上次进度落库时间 + seenVideo map[string]bool // "v:"+去掉扩展名的相对路径 → 远端存在该视频 + seenMeta map[string]bool // "m:"+相对路径 → 远端存在该元数据 + remoteMeta map[string][]remoteMetaItem // "m:"+相对路径 → 远端元数据副本列表(多副本聚合,支持择优比对与冗余清理) seenMetaTarget map[string]cloud.FileEntry seenVideoTarget map[string]cloud.FileEntry activeDownloadPaths map[string]bool // 本地已在排队/进行的下载任务路径(内存去重) @@ -263,9 +269,7 @@ func (s *StrmService) runSync(ctx context.Context, p *model.StrmSyncPath, rec *m syncType: rec.SyncType, seenVideo: map[string]bool{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, - remoteMetaSha1: map[string]string{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, seenVideoTarget: map[string]cloud.FileEntry{}, } @@ -1109,25 +1113,21 @@ func (st *strmSyncState) localSha1Matches(path, remoteSha1 string) bool { return strings.EqualFold(local, remoteSha1) } -// recordRemoteMeta 记录远端存在的元数据索引、文件大小、文件引用及内容 SHA1。 +// recordRemoteMeta 记录远端存在的元数据索引、文件大小、文件引用及内容 SHA1(多副本聚合追加)。 func (st *strmSyncState) recordRemoteMeta(entry cloud.FileEntry, rel string) { st.mu.Lock() + defer st.mu.Unlock() if st.remoteMeta == nil { - st.remoteMeta = map[string]int64{} + st.remoteMeta = map[string][]remoteMetaItem{} } - if st.remoteMetaRef == nil { - st.remoteMetaRef = map[string]string{} - } - if st.remoteMetaSha1 == nil { - st.remoteMetaSha1 = map[string]string{} - } - st.seenMeta["m:"+rel] = true - st.remoteMeta["m:"+rel] = entry.Size - st.remoteMetaRef["m:"+rel] = entry.ID - if sha := usableSha1(entry.Sha1); sha != "" { - st.remoteMetaSha1["m:"+rel] = sha - } - st.mu.Unlock() + key := "m:" + rel + st.seenMeta[key] = true + st.remoteMeta[key] = append(st.remoteMeta[key], remoteMetaItem{ + ID: entry.ID, + Size: entry.Size, + Sha1: usableSha1(entry.Sha1), + MTime: entry.MTime, + }) } func (st *strmSyncState) flushPendingDownloads() { @@ -1379,22 +1379,41 @@ func (st *strmSyncState) scanLocalMetaForUpload() error { return nil } st.mu.Lock() - remoteSize, exists := st.remoteMeta["m:"+rel] - remoteRef := st.remoteMetaRef["m:"+rel] - remoteSha1 := st.remoteMetaSha1["m:"+rel] + entries, exists := st.remoteMeta["m:"+rel] st.mu.Unlock() - if exists && remoteSize == info.Size() { - // 网盘端同名同大小:候选同一文件。115 提供远端 SHA1 时做内容级比对, - // 识别"同大小不同内容"(如 nfo 改一个字符长度不变)避免漏传; - // 无哈希(其他网盘/列表未返回)视为同一文件,保持大小比对兜底。 - if remoteSha1 == "" || st.localSha1Matches(path, remoteSha1) { + + if exists && len(entries) > 0 { + // 择优比对:只要远端存在任一副本的大小匹配且 SHA1 匹配(或无 SHA1),即判定远端已有最新副本,无需上传 + matchedIdx := -1 + for i, it := range entries { + if it.Size == info.Size() { + if it.Sha1 == "" || st.localSha1Matches(path, it.Sha1) { + matchedIdx = i + break + } + } + } + if matchedIdx >= 0 { + // 远端已存在完全一致的副本,跳过上传! + // 若远端还存在其他同名脏副本(副本总数 > 1),在 115 下顺手异步清理其余冗余副本 + if len(entries) > 1 && st.p.Provider == model.StrmProvider115 { + var redundantIDs []string + for i, it := range entries { + if i != matchedIdx && it.ID != "" { + redundantIDs = append(redundantIDs, it.ID) + } + } + if len(redundantIDs) > 0 { + st.cleanupRedundantRemoteFiles(st.uploadRemoteTarget(rel), redundantIDs) + } + } return nil } } - // 入队上传(以本地为准):网盘端不存在;或同名但大小不同(必然内容不同); - // 或同名同大小但 SHA1 不同(精确比对发现的同大小不同内容)。 - // 115 的上传接口不保证同名覆盖,任务携带远端旧文件 ID(RemoteRef), - // 由上传端先删旧文件再上传;WebDAV/OpenList 的 PutFile 本身即覆盖上传。 + + // 入队上传(以本地为准):网盘端不存在;或所有远端副本均内容不同。 + // 115 的上传接口不保证同名覆盖,任务携带远端所有旧副本 ID(RemoteRef,以逗号连接), + // 由上传端一次性批量删除所有旧文件后再上传,彻底根除同名文件堆积。 st.mu.Lock() if st.activeUploadPaths != nil && st.activeUploadPaths[path] { st.mu.Unlock() @@ -1415,8 +1434,14 @@ func (st *strmSyncState) scanLocalMetaForUpload() error { Size: info.Size(), Status: model.StrmTaskPending, } - if exists && st.p.Provider == model.StrmProvider115 { - task.RemoteRef = remoteRef + if exists && len(entries) > 0 && st.p.Provider == model.StrmProvider115 { + var oldIDs []string + for _, it := range entries { + if it.ID != "" { + oldIDs = append(oldIDs, it.ID) + } + } + task.RemoteRef = strings.Join(oldIDs, ",") } st.mu.Lock() st.pendingUploads = append(st.pendingUploads, task) @@ -1431,6 +1456,31 @@ func (st *strmSyncState) scanLocalMetaForUpload() error { }) } +// cleanupRedundantRemoteFiles 异步清理 115 远端某个父目录下的冗余同名旧元数据副本。 +func (st *strmSyncState) cleanupRedundantRemoteFiles(parentCID string, fileIDs []string) { + if len(fileIDs) == 0 || st.provider == nil { + return + } + open115, ok := st.provider.(cloud.OpenAPI115Provider) + if !ok { + return + } + go func() { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + if err := open115.OpenClient().DeleteFiles(ctx, parentCID, fileIDs...); err != nil { + st.s.log.Warn("清理 115 冗余旧元数据副本失败", + zap.String("parent_cid", parentCID), + zap.Strings("file_ids", fileIDs), + zap.Error(err)) + } else { + st.s.log.Info("已清理 115 冗余旧元数据副本", + zap.String("parent_cid", parentCID), + zap.Strings("file_ids", fileIDs)) + } + }() +} + // remoteUploadPath 远端元数据目标路径 = 同步目录远端根 + 相对路径。 func (st *strmSyncState) remoteUploadPath(rel string) string { root := strings.TrimRight(normalizeRemotePath(st.p.RemotePath), "/") diff --git a/internal/service/strm_sync_test.go b/internal/service/strm_sync_test.go index 0d5ac76..214243f 100644 --- a/internal/service/strm_sync_test.go +++ b/internal/service/strm_sync_test.go @@ -368,19 +368,18 @@ func TestScanLocalMetaForUpload(t *testing.T) { } st := &strmSyncState{ - s: svc, - ctx: context.Background(), - p: p, - cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"jpg", "nfo"}}, - rec: &model.StrmSyncRecord{}, - seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, + s: svc, + ctx: context.Background(), + p: p, + cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"jpg", "nfo"}}, + rec: &model.StrmSyncRecord{}, + seenMeta: map[string]bool{}, + remoteMeta: map[string][]remoteMetaItem{}, } // 模拟远端已存在 poster.jpg(与本地同一文件)和 tvshow.nfo(与本地不同) - st.remoteMeta["m:动漫/poster.jpg"] = int64(len("poster-data")) - st.remoteMeta["m:动漫/tvshow.nfo"] = 999 + st.remoteMeta["m:动漫/poster.jpg"] = []remoteMetaItem{{ID: "f1", Size: int64(len("poster-data"))}} + st.remoteMeta["m:动漫/tvshow.nfo"] = []remoteMetaItem{{ID: "f2", Size: 999}} if err := st.scanLocalMetaForUpload(); err != nil { t.Fatalf("scanLocalMetaForUpload failed: %v", err) @@ -442,17 +441,15 @@ func TestScanLocalMetaForUpload115CarriesRemoteRef(t *testing.T) { } st := &strmSyncState{ - s: svc, - ctx: context.Background(), - p: p, - cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"jpg", "nfo"}}, - rec: &model.StrmSyncRecord{}, - seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, + s: svc, + ctx: context.Background(), + p: p, + cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"jpg", "nfo"}}, + rec: &model.StrmSyncRecord{}, + seenMeta: map[string]bool{}, + remoteMeta: map[string][]remoteMetaItem{}, } - st.remoteMeta["m:movie.nfo"] = 1 - st.remoteMetaRef["m:movie.nfo"] = "file-42" + st.remoteMeta["m:movie.nfo"] = []remoteMetaItem{{ID: "file-42", Size: 1}} if err := st.scanLocalMetaForUpload(); err != nil { t.Fatalf("scanLocalMetaForUpload failed: %v", err) @@ -522,7 +519,7 @@ func TestPruneLocalKeepsLocalMeta(t *testing.T) { syncType: model.StrmSyncTypeFull, seenVideo: map[string]bool{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, + remoteMeta: map[string][]remoteMetaItem{}, } // 本次远端扫描既没有看到视频,也没有看到任何元数据 if err := st.pruneLocal(); err != nil { @@ -569,8 +566,7 @@ func TestHandleMetaKeepsLocalWhenUploadEnabled(t *testing.T) { cfg: &strmPathConfig{DownloadMeta: true, UploadMeta: true, MetaExt: []string{"nfo"}}, rec: &model.StrmSyncRecord{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, } @@ -630,22 +626,21 @@ func TestScanLocalMetaForUploadSha1Identity(t *testing.T) { UploadMeta: true, } st := &strmSyncState{ - s: svc, - ctx: context.Background(), - p: p, - cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"nfo"}}, - rec: &model.StrmSyncRecord{}, - seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, - remoteMetaSha1: map[string]string{}, + s: svc, + ctx: context.Background(), + p: p, + cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"nfo"}}, + rec: &model.StrmSyncRecord{}, + seenMeta: map[string]bool{}, + remoteMeta: map[string][]remoteMetaItem{}, } // 115 返回大写 SHA1,本地计算为小写:同时验证大小写不敏感比对 - st.remoteMeta["m:same.nfo"] = int64(len("same-content")) - st.remoteMetaSha1["m:same.nfo"] = strings.ToUpper(sameSha) - st.remoteMeta["m:diff.nfo"] = int64(len("diff-content")) - st.remoteMetaSha1["m:diff.nfo"] = strings.ToUpper(otherSha) - st.remoteMetaRef["m:diff.nfo"] = "old-diff-1" + st.remoteMeta["m:same.nfo"] = []remoteMetaItem{ + {ID: "same-1", Size: int64(len("same-content")), Sha1: strings.ToUpper(sameSha)}, + } + st.remoteMeta["m:diff.nfo"] = []remoteMetaItem{ + {ID: "old-diff-1", Size: int64(len("diff-content")), Sha1: strings.ToUpper(otherSha)}, + } if err := st.scanLocalMetaForUpload(); err != nil { t.Fatalf("scanLocalMetaForUpload failed: %v", err) @@ -703,9 +698,7 @@ func TestHandleMetaSha1Identity(t *testing.T) { cfg: &strmPathConfig{DownloadMeta: true, UploadMeta: false, MetaExt: []string{"nfo"}}, rec: &model.StrmSyncRecord{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, - remoteMetaSha1: map[string]string{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, } @@ -802,7 +795,7 @@ func TestWalkRemoteConcurrent(t *testing.T) { rec: &model.StrmSyncRecord{}, seenVideo: map[string]bool{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, + remoteMeta: map[string][]remoteMetaItem{}, } if err := st.walkRemote(); err != nil { t.Fatalf("walkRemote failed: %v", err) @@ -911,7 +904,7 @@ func TestStrmDuplicateFileConflictResolution(t *testing.T) { cfg: &strmPathConfig{DownloadMeta: true, MetaExt: []string{"nfo"}}, rec: &model.StrmSyncRecord{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, seenVideoTarget: map[string]cloud.FileEntry{}, } @@ -950,7 +943,7 @@ func TestStrmDuplicateFileConflictResolution(t *testing.T) { cfg: &strmPathConfig{DownloadMeta: true, MetaExt: []string{"nfo"}}, rec: &model.StrmSyncRecord{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, seenVideoTarget: map[string]cloud.FileEntry{}, } @@ -1019,7 +1012,7 @@ func TestWalk115FlatAbortsOnDirResolveFailure(t *testing.T) { dirCache: sync.Map{}, seenVideo: map[string]bool{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, + remoteMeta: map[string][]remoteMetaItem{}, } err := st.walk115Flat(oc) if err == nil { @@ -1088,9 +1081,7 @@ func TestWalk115FlatConcurrentProcessing(t *testing.T) { dirCache: sync.Map{}, seenVideo: map[string]bool{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, - remoteMetaSha1: map[string]string{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, seenVideoTarget: map[string]cloud.FileEntry{}, } @@ -1145,9 +1136,7 @@ func TestWalk115FlatConcurrentProcessing(t *testing.T) { dirCache: sync.Map{}, seenVideo: map[string]bool{}, seenMeta: map[string]bool{}, - remoteMeta: map[string]int64{}, - remoteMetaRef: map[string]string{}, - remoteMetaSha1: map[string]string{}, + remoteMeta: map[string][]remoteMetaItem{}, seenMetaTarget: map[string]cloud.FileEntry{}, seenVideoTarget: map[string]cloud.FileEntry{}, } @@ -1206,3 +1195,127 @@ func TestCloud115FullPath(t *testing.T) { t.Fatalf("root full path = %q, want empty", got) } } + +// TestScanLocalMetaForUploadMultipleCopiesBestEffortMatch 验证: +// 当远端同一路径存在多个副本(1个与本地一致的副本 + 1个脏副本)时, +// 能够择优识别出匹配的副本,跳过上传(uploaded = 0),避免盲盒覆盖导致的重复上传。 +func TestScanLocalMetaForUploadMultipleCopiesBestEffortMatch(t *testing.T) { + svc := testStrmService(t) + localDir := t.TempDir() + + nfoPath := filepath.Join(localDir, "test.nfo") + writeFile(t, nfoPath, "correct-content") + correctSha, err := cloud115.FileSHA1(nfoPath) + if err != nil { + t.Fatal(err) + } + + p := &model.StrmSyncPath{ + Base: model.Base{ID: "multi-copy-match-path"}, + Provider: model.StrmProvider115, + RemotePath: "root-cid", + LocalPath: localDir, + UploadMeta: true, + } + + st := &strmSyncState{ + s: svc, + ctx: context.Background(), + p: p, + cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"nfo"}}, + rec: &model.StrmSyncRecord{}, + seenMeta: map[string]bool{}, + remoteMeta: map[string][]remoteMetaItem{}, + } + + // 模拟远端存在两个同名副本:一个脏副本(较早),一个正确副本(较晚) + st.recordRemoteMeta(cloud.FileEntry{ + ID: "stale-id", + Name: "test.nfo", + Size: int64(len("stale-dirty-content")), + Sha1: "STALE_SHA1", + MTime: 1000, + }, "test.nfo") + + st.recordRemoteMeta(cloud.FileEntry{ + ID: "correct-id", + Name: "test.nfo", + Size: int64(len("correct-content")), + Sha1: correctSha, + MTime: 2000, + }, "test.nfo") + + if err := st.scanLocalMetaForUpload(); err != nil { + t.Fatalf("scanLocalMetaForUpload failed: %v", err) + } + + // 择优匹配:命中 correct-id,不产生上传任务 + tasks, _, err := svc.repo.StrmUpload.List(context.Background(), "", 1, 10) + if err != nil { + t.Fatal(err) + } + if len(tasks) != 0 { + t.Fatalf("expected 0 upload tasks when a matching copy exists, got %d: %v", len(tasks), taskNames(tasks)) + } +} + +// TestScanLocalMetaForUploadMultipleCopiesAllStale 验证: +// 当远端同一路径存在多个副本,且所有副本均与本地内容不一致时, +// 上传任务应携带所有旧副本的 ID(逗号分隔),以便上传前批量清理所有旧副本。 +func TestScanLocalMetaForUploadMultipleCopiesAllStale(t *testing.T) { + svc := testStrmService(t) + localDir := t.TempDir() + + nfoPath := filepath.Join(localDir, "test.nfo") + writeFile(t, nfoPath, "brand-new-content") + + p := &model.StrmSyncPath{ + Base: model.Base{ID: "multi-copy-stale-path"}, + Provider: model.StrmProvider115, + RemotePath: "root-cid", + LocalPath: localDir, + UploadMeta: true, + } + + st := &strmSyncState{ + s: svc, + ctx: context.Background(), + p: p, + cfg: &strmPathConfig{UploadMeta: true, MetaExt: []string{"nfo"}}, + rec: &model.StrmSyncRecord{}, + seenMeta: map[string]bool{}, + remoteMeta: map[string][]remoteMetaItem{}, + } + + // 模拟远端存在两个不同大小和哈希的旧副本 + st.recordRemoteMeta(cloud.FileEntry{ + ID: "old-1", + Name: "test.nfo", + Size: 10, + Sha1: "OLD_SHA1", + MTime: 1000, + }, "test.nfo") + + st.recordRemoteMeta(cloud.FileEntry{ + ID: "old-2", + Name: "test.nfo", + Size: 20, + Sha1: "OLD_SHA2", + MTime: 2000, + }, "test.nfo") + + if err := st.scanLocalMetaForUpload(); err != nil { + t.Fatalf("scanLocalMetaForUpload failed: %v", err) + } + + tasks, _, err := svc.repo.StrmUpload.List(context.Background(), "", 1, 10) + if err != nil { + t.Fatal(err) + } + if len(tasks) != 1 { + t.Fatalf("expected 1 upload task, got %d", len(tasks)) + } + if tasks[0].RemoteRef != "old-1,old-2" { + t.Fatalf("expected RemoteRef to be 'old-1,old-2', got %q", tasks[0].RemoteRef) + } +}