From 0384017e98e5e0dda8367d11567f2d5842b7bead Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Wed, 26 Aug 2026 14:50:06 +0800 Subject: [PATCH] 6 --- internal/service/strm_sync.go | 52 +++++++++++++++----- internal/service/strm_sync_test.go | 78 +++++++++++++++++++++++++++++- 2 files changed, 117 insertions(+), 13 deletions(-) diff --git a/internal/service/strm_sync.go b/internal/service/strm_sync.go index 0982c7c..f16a1b6 100644 --- a/internal/service/strm_sync.go +++ b/internal/service/strm_sync.go @@ -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 diff --git a/internal/service/strm_sync_test.go b/internal/service/strm_sync_test.go index c5254a7..8d87a8d 100644 --- a/internal/service/strm_sync_test.go +++ b/internal/service/strm_sync_test.go @@ -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) } } +