mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-28 03:06:38 +08:00
Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a71a18ce82 | |||
| 5a189a44fc | |||
| fb84c62e9a | |||
| 355fd06036 | |||
| 4173caac5d |
@@ -87,15 +87,15 @@ func (p *cloudDrive2Provider) Resolve(ctx context.Context, fileRef string) (*Dir
|
||||
if ref == "/" {
|
||||
return nil, fmt.Errorf("%s: file reference required", p.name)
|
||||
}
|
||||
if p.typ == TypeOpenList && isCloudVideoPlaybackCandidate(ref) {
|
||||
if p.apiBase == nil {
|
||||
return nil, fmt.Errorf("%s: pure 302 playback requires an OpenList API server address; configure server/api_url so /api/fs/get can return raw_url", p.name)
|
||||
}
|
||||
if p.typ == TypeOpenList && p.apiBase != nil {
|
||||
link, err := p.resolveOpenListAPIDirect(ctx, ref)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s: pure 302 playback requires OpenList raw_url for %s: %w", p.name, ref, err)
|
||||
if err == nil {
|
||||
return link, nil
|
||||
}
|
||||
// API 获取直链失败:非视频文件(元数据)回退到 WebDAV;视频文件报错
|
||||
if isCloudVideoPlaybackCandidate(ref) {
|
||||
return nil, fmt.Errorf("%s: resolve download URL for %s via API failed: %w", p.name, ref, err)
|
||||
}
|
||||
return link, nil
|
||||
}
|
||||
if p.typ == TypeCloudDrive2 && isCloudVideoPlaybackCandidate(ref) {
|
||||
link, err := p.resolveCloudDAVRedirectDirect(ctx, ref)
|
||||
|
||||
@@ -61,9 +61,10 @@ func TestOpenListWebDAVListAndResolve(t *testing.T) {
|
||||
if len(entries) != 1 || entries[0].ID != "/Cloud/Movie.mkv" || entries[0].Size != 1024 {
|
||||
t.Fatalf("entries = %#v", entries)
|
||||
}
|
||||
// Video file: API fails → error (no WebDAV fallback for video)
|
||||
_, err = p.Resolve(context.Background(), entries[0].ID)
|
||||
if err == nil || !strings.Contains(err.Error(), "pure 302 playback requires OpenList raw_url") {
|
||||
t.Fatalf("openlist video resolve should require raw_url instead of WebDAV proxy fallback, err=%v", err)
|
||||
if err == nil || !strings.Contains(err.Error(), "resolve download URL") || !strings.Contains(err.Error(), "via API failed") {
|
||||
t.Fatalf("openlist video resolve should error on API failure, err=%v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -195,10 +195,43 @@ func TestOpenListResolveDoesNotFallbackToWebDAVWhenAPIRawURLFails(t *testing.T)
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err = p.Resolve(context.Background(), "/Cloud/Movie.mkv")
|
||||
if err == nil || !strings.Contains(err.Error(), "pure 302 playback requires OpenList raw_url") {
|
||||
t.Fatalf("resolve error = %v, want raw_url requirement", err)
|
||||
if err == nil || !strings.Contains(err.Error(), "resolve download URL") || !strings.Contains(err.Error(), "via API failed") {
|
||||
t.Fatalf("resolve error = %v, want API resolve failure", err)
|
||||
}
|
||||
if davSeen {
|
||||
t.Fatal("openlist video resolve fell back to WebDAV after raw_url failure")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenListResolveMetadataUsesAPIInsteadOfWebDAV(t *testing.T) {
|
||||
var gotPath, gotAuth string
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
gotPath = r.URL.Path
|
||||
gotAuth = r.Header.Get("Authorization")
|
||||
if r.Method != http.MethodPost || r.URL.Path != "/api/fs/get" {
|
||||
t.Fatalf("unexpected request %s %s; metadata should use API, not WebDAV", r.Method, r.URL.Path)
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_, _ = w.Write([]byte(`{"code":200,"data":{"raw_url":"https://cdn.example.test/poster.jpg?sign=1"}}`))
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
p, err := New(TypeOpenList, map[string]any{"server": srv.URL, "token": "alist-token"}, srv.Client())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// .nfo metadata file should use API, not WebDAV
|
||||
link, err := p.Resolve(context.Background(), "/Cloud/Movie/Movie.nfo")
|
||||
if err != nil {
|
||||
t.Fatalf("resolve: %v", err)
|
||||
}
|
||||
if gotPath != "/api/fs/get" {
|
||||
t.Fatalf("api path = %q, want /api/fs/get (metadata should not use WebDAV)", gotPath)
|
||||
}
|
||||
if gotAuth != "alist-token" {
|
||||
t.Fatalf("Authorization = %q, want token", gotAuth)
|
||||
}
|
||||
if link.URL != "https://cdn.example.test/poster.jpg?sign=1" {
|
||||
t.Fatalf("url = %q", link.URL)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -366,9 +366,11 @@ func (s *StrmService) processUpload115(ctx context.Context, task *model.StrmUplo
|
||||
finish(model.StrmTaskFailed, "该网盘不支持元数据上传")
|
||||
return
|
||||
}
|
||||
// 以本地为准:网盘端已有同名但内容不同的旧元数据时,先批量删除所有旧副本再上传。
|
||||
// 115 的上传接口不保证同名覆盖,直接上传可能产生同名重复文件;删除失败则
|
||||
// 任务重试(旧文件 ID 失效的场景会在下次同步后自动修复)。
|
||||
// 以本地为准:网盘端已有同名但内容不同的旧元数据时,先尝试批量删除所有旧副本再上传。
|
||||
// 115 的上传接口不保证同名覆盖,直接上传可能产生同名重复文件。
|
||||
// 删除失败时不中止任务——继续上传新文件,旧副本交由下次同步的 cleanupBatchRedundantFiles
|
||||
// 按目录批量清理(下次同步会看到新旧两个版本,命中新版本后把旧版本 cid 收入 pendingDeletes
|
||||
// 异步删除)。这样避免了「删旧失败 → 任务重试 → 再次删旧失败 → 永远无法上传」的死循环。
|
||||
if task.RemoteRef != "" {
|
||||
open115, ok := provider.(cloud.OpenAPI115Provider)
|
||||
if !ok {
|
||||
@@ -377,8 +379,11 @@ func (s *StrmService) processUpload115(ctx context.Context, task *model.StrmUplo
|
||||
}
|
||||
refs := strings.Split(task.RemoteRef, ",")
|
||||
if err := open115.OpenClient().DeleteFiles(ctx, task.RemotePath, refs...); err != nil {
|
||||
s.uploadTaskFailWithRetry(task, "删除网盘旧元数据失败:"+err.Error())
|
||||
return
|
||||
s.log.Warn("删除网盘旧元数据失败,跳过删除继续上传新文件",
|
||||
zap.String("task_id", task.ID),
|
||||
zap.String("local_path", task.LocalPath),
|
||||
zap.Error(err))
|
||||
// 不 return:继续上传新文件,旧副本由下次同步清理
|
||||
}
|
||||
}
|
||||
f, err := os.Open(task.LocalPath)
|
||||
|
||||
@@ -25,6 +25,8 @@ import (
|
||||
"github.com/truewhile/MeBox/internal/service/cloud115"
|
||||
)
|
||||
|
||||
var errFallbackToWalkRemote = errors.New("fallback to walk remote")
|
||||
|
||||
// remoteMetaItem 记录远端存在的单个元数据文件副本信息(大小、文件ID、内容SHA1、修改时间)。
|
||||
type remoteMetaItem struct {
|
||||
ID string
|
||||
@@ -352,8 +354,17 @@ func (st *strmSyncState) run() error {
|
||||
|
||||
if st.provider != nil {
|
||||
if open115, ok := st.provider.(cloud.OpenAPI115Provider); ok && st.p.Provider == model.StrmProvider115 {
|
||||
if err := st.walk115Flat(open115.OpenClient()); err != nil {
|
||||
return err
|
||||
err := st.walk115Flat(open115.OpenClient())
|
||||
if err != nil {
|
||||
if errors.Is(err, errFallbackToWalkRemote) {
|
||||
st.s.log.Warn("115: 扁平列表文件数超限,自动降级为传统并发递归同步",
|
||||
zap.String("path_id", st.p.ID))
|
||||
if errWalk := st.walkRemote(); errWalk != nil {
|
||||
return errWalk
|
||||
}
|
||||
} else {
|
||||
return err
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if err := st.walkRemote(); err != nil {
|
||||
@@ -675,6 +686,13 @@ func (st *strmSyncState) walk115Flat(open115 *cloud115.OpenClient) error {
|
||||
return fmt.Errorf("115: 获取文件列表失败:%w", err)
|
||||
}
|
||||
|
||||
// 115 的扁平化列表(搜索底层)对 offset + limit 有 10000 的最大深度限制。
|
||||
// 当扁平模式下的总文件数 >= 9500 时,强行拒绝继续扁平拉取,而是抛出降级错误,
|
||||
// 让外层回退到使用普通的按目录并发递归(walkRemote),以免截断导致后排文件被误删/重传。
|
||||
if totalCount >= 9500 {
|
||||
return errFallbackToWalkRemote
|
||||
}
|
||||
|
||||
st.updateSyncMessage(fmt.Sprintf("正在拉取远端文件列表 (共 %d 个文件)...", totalCount))
|
||||
|
||||
allFiles := make([]cloud115.RemoteFile, 0, totalCount)
|
||||
@@ -1349,8 +1367,12 @@ func (st *strmSyncState) scanLocalMetaForUpload() error {
|
||||
st.activeUploadPaths = map[string]bool{}
|
||||
}
|
||||
}
|
||||
var pendingDeletes map[string][]string
|
||||
if st.p.Provider == model.StrmProvider115 {
|
||||
pendingDeletes = map[string][]string{}
|
||||
}
|
||||
localRoot := filepath.Clean(st.p.LocalPath)
|
||||
return filepath.WalkDir(localRoot, func(path string, d os.DirEntry, err error) error {
|
||||
err := filepath.WalkDir(localRoot, func(path string, d os.DirEntry, err error) error {
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
@@ -1395,22 +1417,27 @@ func (st *strmSyncState) scanLocalMetaForUpload() error {
|
||||
}
|
||||
if matchedIdx >= 0 {
|
||||
// 远端已存在完全一致的副本,跳过上传!
|
||||
// 若远端还存在其他同名脏副本(副本总数 > 1),在 115 下顺手异步清理其余冗余副本
|
||||
// 若远端还存在其他同名脏副本(副本总数 > 1),在 115 下收集待删除 ID,稍后按目录批量删除
|
||||
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)
|
||||
parentCID := st.uploadRemoteTarget(rel)
|
||||
if parentCID != "" {
|
||||
for i, it := range entries {
|
||||
if i != matchedIdx && it.ID != "" {
|
||||
pendingDeletes[parentCID] = append(pendingDeletes[parentCID], it.ID)
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(redundantIDs) > 0 {
|
||||
st.cleanupRedundantRemoteFiles(st.uploadRemoteTarget(rel), redundantIDs)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
remoteTarget := st.uploadRemoteTarget(rel)
|
||||
if st.p.Provider == model.StrmProvider115 && remoteTarget == "" {
|
||||
// 115 远端不存在对应父目录(如孤儿子目录),禁止降级到根目录上传以防错位死循环
|
||||
return nil
|
||||
}
|
||||
|
||||
// 入队上传(以本地为准):网盘端不存在;或所有远端副本均内容不同。
|
||||
// 115 的上传接口不保证同名覆盖,任务携带远端所有旧副本 ID(RemoteRef,以逗号连接),
|
||||
// 由上传端一次性批量删除所有旧文件后再上传,彻底根除同名文件堆积。
|
||||
@@ -1430,7 +1457,7 @@ func (st *strmSyncState) scanLocalMetaForUpload() error {
|
||||
Provider: st.p.Provider,
|
||||
FileName: filepath.Base(rel),
|
||||
LocalPath: path,
|
||||
RemotePath: st.uploadRemoteTarget(rel),
|
||||
RemotePath: remoteTarget,
|
||||
Size: info.Size(),
|
||||
Status: model.StrmTaskPending,
|
||||
}
|
||||
@@ -1454,11 +1481,20 @@ func (st *strmSyncState) scanLocalMetaForUpload() error {
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// 批量清理 115 冗余旧元数据副本(按父目录聚合,一次请求批量删除该目录下全部冗余文件,彻底避免并发触发限流器超时)
|
||||
if len(pendingDeletes) > 0 {
|
||||
st.cleanupBatchRedundantFiles(pendingDeletes)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// cleanupRedundantRemoteFiles 异步清理 115 远端某个父目录下的冗余同名旧元数据副本。
|
||||
func (st *strmSyncState) cleanupRedundantRemoteFiles(parentCID string, fileIDs []string) {
|
||||
if len(fileIDs) == 0 || st.provider == nil {
|
||||
// cleanupBatchRedundantFiles 异步按父目录批量清理 115 远端冗余旧元数据副本。
|
||||
func (st *strmSyncState) cleanupBatchRedundantFiles(deletesByParent map[string][]string) {
|
||||
if len(deletesByParent) == 0 || st.provider == nil {
|
||||
return
|
||||
}
|
||||
open115, ok := st.provider.(cloud.OpenAPI115Provider)
|
||||
@@ -1466,17 +1502,28 @@ func (st *strmSyncState) cleanupRedundantRemoteFiles(parentCID string, fileIDs [
|
||||
return
|
||||
}
|
||||
go func() {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
|
||||
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))
|
||||
totalPruned := 0
|
||||
for parentCID, fileIDs := range deletesByParent {
|
||||
if ctx.Err() != nil {
|
||||
break
|
||||
}
|
||||
if len(fileIDs) == 0 {
|
||||
continue
|
||||
}
|
||||
if err := open115.OpenClient().DeleteFiles(ctx, parentCID, fileIDs...); err != nil {
|
||||
st.s.log.Warn("批量清理 115 冗余旧元数据副本失败",
|
||||
zap.String("parent_cid", parentCID),
|
||||
zap.Int("count", len(fileIDs)),
|
||||
zap.Error(err))
|
||||
} else {
|
||||
totalPruned += len(fileIDs)
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
}
|
||||
if totalPruned > 0 {
|
||||
st.s.log.Info("已完成批量清理 115 冗余旧元数据副本", zap.Int("total_pruned", totalPruned))
|
||||
}
|
||||
}()
|
||||
}
|
||||
@@ -1508,9 +1555,9 @@ func (st *strmSyncState) uploadRemoteTarget(rel string) string {
|
||||
if cid, ok := st.dirPathToID[dir]; ok && cid != "" {
|
||||
return cid
|
||||
}
|
||||
// 父目录未在缓存中(父目录可能本次未扫描到),降级为用户配置的同步根 cid,
|
||||
// 由上传端尽力处理(可能失败记日志,不影响下载)。
|
||||
return st.p.RemotePath
|
||||
// 子目录在 115 远端不存在(本地孤儿子目录或远端已删除该分类文件夹),
|
||||
// 返回空串,禁止降级回退到根目录上传以防污染根目录与错位死循环。
|
||||
return ""
|
||||
}
|
||||
return st.remoteUploadPath(rel)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user