mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-28 11:16:37 +08:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 22b7290ee1 | |||
| 14037d5dea | |||
| 152db3fb9f | |||
| 0413d123da | |||
| 5f6bd7b5cd | |||
| 3c25bb5d61 |
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -5,6 +5,8 @@ package cloud115
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -181,6 +183,24 @@ func (u *OSSMultipartUploader) UploadFileWithResult(ctx context.Context, input O
|
||||
return completeParts[i].PartNumber < completeParts[j].PartNumber
|
||||
})
|
||||
|
||||
// 115 下发的 callback / callback_var 是 JSON 字符串,而 OSS CompleteMultipartUpload
|
||||
// 要求 callback 参数为 Base64 编码后的 JSON,否则报 "The callback configuration is
|
||||
// not base64 encoded"。这里把两者转为 Base64 后再提交(参考 QMediaSync 的
|
||||
// BuildOSSCallbackHeaders)。
|
||||
cb := input.Callback
|
||||
cbVar := input.CallbackVar
|
||||
if cb == "" {
|
||||
return OSSMultipartUploadResult{}, errors.New("OSS callback 为空")
|
||||
}
|
||||
if !json.Valid([]byte(cb)) {
|
||||
return OSSMultipartUploadResult{}, errors.New("解析 callback 失败:不是合法 JSON")
|
||||
}
|
||||
if cbVar == "" {
|
||||
cbVar = "{}"
|
||||
}
|
||||
if !json.Valid([]byte(cbVar)) {
|
||||
return OSSMultipartUploadResult{}, errors.New("解析 callback_var 失败:不是合法 JSON")
|
||||
}
|
||||
completeResult, err := u.client.CompleteMultipartUpload(ctx, &oss.CompleteMultipartUploadRequest{
|
||||
Bucket: oss.Ptr(input.Bucket),
|
||||
Key: oss.Ptr(input.Object),
|
||||
@@ -188,8 +208,8 @@ func (u *OSSMultipartUploader) UploadFileWithResult(ctx context.Context, input O
|
||||
CompleteMultipartUpload: &oss.CompleteMultipartUpload{
|
||||
Parts: completeParts,
|
||||
},
|
||||
Callback: oss.Ptr(input.Callback),
|
||||
CallbackVar: oss.Ptr(input.CallbackVar),
|
||||
Callback: oss.Ptr(base64.StdEncoding.EncodeToString([]byte(cb))),
|
||||
CallbackVar: oss.Ptr(base64.StdEncoding.EncodeToString([]byte(cbVar))),
|
||||
})
|
||||
if err != nil {
|
||||
return OSSMultipartUploadResult{}, fmt.Errorf("完成 OSS multipart 失败:%w", err)
|
||||
|
||||
@@ -40,7 +40,9 @@ func FileSHA1Partial(path string, start, end int64) (string, error) {
|
||||
}
|
||||
length := end - start + 1
|
||||
h := sha1.New()
|
||||
if _, err := io.CopyN(h, f, length); err != nil {
|
||||
// io.CopyN 在文件不足 length 字节时会返回 io.EOF,导致小文件(如小于 128 KiB 的
|
||||
// 元数据图片)无法上传。这里只拷贝实际读到的字节,文件尾对齐到区间终点即可。
|
||||
if _, err := io.CopyN(h, f, length); err != nil && err != io.EOF {
|
||||
return "", err
|
||||
}
|
||||
return hex.EncodeToString(h.Sum(nil)), nil
|
||||
|
||||
@@ -1,9 +1,14 @@
|
||||
package cloud115
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/aliyun/alibabacloud-oss-go-sdk-v2/oss"
|
||||
)
|
||||
|
||||
func TestFileSHA1(t *testing.T) {
|
||||
@@ -39,6 +44,26 @@ func TestFileSHA1Partial(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestFileSHA1PartialSmallerThanWindow 回归测试:经典 bug 是 io.CopyN 在文件不足
|
||||
// length 字节时返回 io.EOF。115 上传固定用 [0,128*1024-1] 窗口计算 preid,导致所有
|
||||
// 小于 128 KiB 的元数据文件(如海报/缩略图)上传必然失败。
|
||||
func TestFileSHA1PartialSmallerThanWindow(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, "small.bin")
|
||||
// 6 字节小文件,不足 128 KiB 窗口
|
||||
if err := os.WriteFile(path, []byte("abcdef"), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sum, err := FileSHA1Partial(path, 0, 128*1024-1)
|
||||
if err != nil {
|
||||
t.Fatalf("compute partial sha1 for small file should not fail: %v", err)
|
||||
}
|
||||
// 应等于整个文件(6 字节)的 sha1
|
||||
if sum != "1f8ac10f23c5b5bc1167bda84b833e5c057a77d2" {
|
||||
t.Errorf("unexpected partial sha1: %s", sum)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseSignCheckRange(t *testing.T) {
|
||||
rng, err := parseSignCheckRange("0-131071")
|
||||
if err != nil {
|
||||
@@ -92,3 +117,81 @@ func TestBaseNameOf(t *testing.T) {
|
||||
t.Errorf("got %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
// fakeCallbackOSSClient 捕获 CompleteMultipartUpload 收到的 callback / callback_var,
|
||||
// 用于断言已经 Base64 编码(116 要求 callback 必须是 Base64 后的 JSON,否则报
|
||||
// "The callback configuration is not base64 encoded")。
|
||||
type fakeCallbackOSSClient struct {
|
||||
capturedCallback string
|
||||
capturedCallbackVar string
|
||||
}
|
||||
|
||||
func (c *fakeCallbackOSSClient) InitiateMultipartUpload(_ context.Context, _ *oss.InitiateMultipartUploadRequest, _ ...func(*oss.Options)) (*oss.InitiateMultipartUploadResult, error) {
|
||||
return &oss.InitiateMultipartUploadResult{UploadId: oss.Ptr("upload-new")}, nil
|
||||
}
|
||||
func (c *fakeCallbackOSSClient) UploadPart(_ context.Context, r *oss.UploadPartRequest, _ ...func(*oss.Options)) (*oss.UploadPartResult, error) {
|
||||
if r.Body != nil {
|
||||
_, _ = io.Copy(io.Discard, r.Body)
|
||||
}
|
||||
return &oss.UploadPartResult{ETag: oss.Ptr("etag-1")}, nil
|
||||
}
|
||||
func (c *fakeCallbackOSSClient) ListParts(context.Context, *oss.ListPartsRequest, ...func(*oss.Options)) (*oss.ListPartsResult, error) {
|
||||
return &oss.ListPartsResult{}, nil
|
||||
}
|
||||
func (c *fakeCallbackOSSClient) CompleteMultipartUpload(_ context.Context, r *oss.CompleteMultipartUploadRequest, _ ...func(*oss.Options)) (*oss.CompleteMultipartUploadResult, error) {
|
||||
c.capturedCallback = *r.Callback
|
||||
c.capturedCallbackVar = *r.CallbackVar
|
||||
return &oss.CompleteMultipartUploadResult{
|
||||
CallbackResult: map[string]any{
|
||||
"state": true,
|
||||
"data": map[string]any{"file_id": "file-1", "pick_code": "pick-1"},
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
func (c *fakeCallbackOSSClient) AbortMultipartUpload(context.Context, *oss.AbortMultipartUploadRequest, ...func(*oss.Options)) (*oss.AbortMultipartUploadResult, error) {
|
||||
return &oss.AbortMultipartUploadResult{}, nil
|
||||
}
|
||||
|
||||
// TestCompleteMultipartUploadCallbackBase64 回归测试:OSS CompleteMultipartUpload 的
|
||||
// callback 必须 Base64 编码,否则报 "The callback configuration is not base64 encoded",
|
||||
// 导致大于 128 KiB 的元数据文件上传失败。
|
||||
func TestCompleteMultipartUploadCallbackBase64(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, "big.bin")
|
||||
data := make([]byte, 8) // 8 字节,PartSize=8 → 1 part
|
||||
if err := os.WriteFile(path, data, 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fake := &fakeCallbackOSSClient{}
|
||||
uploader := &OSSMultipartUploader{client: fake}
|
||||
|
||||
callback := `{"callbackUrl":"http://uplb.115.com/3.0/completeupload.php"}`
|
||||
callbackVar := `{"x:pick_code":"abc"}`
|
||||
_, err := uploader.UploadFileWithResult(context.Background(), OSSMultipartUploadInput{
|
||||
Bucket: "bucket-1",
|
||||
Object: "object-1",
|
||||
Callback: callback,
|
||||
CallbackVar: callbackVar,
|
||||
FilePath: path,
|
||||
FileSize: int64(len(data)),
|
||||
PartSize: 8,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("multipart 上传失败:%v", err)
|
||||
}
|
||||
// 捕获的 callback 必须是合法 Base64,且解码后与原 JSON 一致
|
||||
cbBytes, err := base64.StdEncoding.DecodeString(fake.capturedCallback)
|
||||
if err != nil {
|
||||
t.Fatalf("callback 未 Base64 编码:%v (raw=%q)", err, fake.capturedCallback)
|
||||
}
|
||||
if string(cbBytes) != callback {
|
||||
t.Errorf("callback 解码后 = %s,期望 %s", cbBytes, callback)
|
||||
}
|
||||
cbvBytes, err := base64.StdEncoding.DecodeString(fake.capturedCallbackVar)
|
||||
if err != nil {
|
||||
t.Fatalf("callback_var 未 Base64 编码:%v (raw=%q)", err, fake.capturedCallbackVar)
|
||||
}
|
||||
if string(cbvBytes) != callbackVar {
|
||||
t.Errorf("callback_var 解码后 = %s,期望 %s", cbvBytes, callbackVar)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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{
|
||||
|
||||
+208
-133
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user