From f7fec93d44ee88851d5b46fc129f09a15a740f6b Mon Sep 17 00:00:00 2001 From: truewhile <62226914+truewhile@users.noreply.github.com> Date: Sun, 6 Sep 2026 17:05:56 +0800 Subject: [PATCH] fix(upload): eliminate concurrent temp file name collision and support direct local file upload for 115 --- internal/service/cloud/pan115_openapi.go | 46 +++++++++++++---------- internal/service/strm_queue.go | 48 +++++++++++++++--------- 2 files changed, 56 insertions(+), 38 deletions(-) diff --git a/internal/service/cloud/pan115_openapi.go b/internal/service/cloud/pan115_openapi.go index 749cfb7..6fe00ef 100644 --- a/internal/service/cloud/pan115_openapi.go +++ b/internal/service/cloud/pan115_openapi.go @@ -123,35 +123,41 @@ func (p *openAPI115Provider) ResolveBatch(ctx context.Context, fileRefs []string // OpenClient 暴露底层客户端(token 刷新用)。 func (p *openAPI115Provider) OpenClient() *cloud115.OpenClient { return p.c } +// PutLocalFile 直接上传本地文件,避免通过 io.Reader 复制临时文件产生的磁盘开销与并发重命名碰撞。 +func (p *openAPI115Provider) PutLocalFile(ctx context.Context, parentCID, localPath string) error { + _, err := p.c.Upload(ctx, localPath, parentCID, "", "") + return err +} + // PutFileNamed 把本地元数据上传到 115 指定父目录(parentCID 为父目录 cid)。 -// io.Reader 无法携带文件名,因此走独立的 named 上传接口。将内容落为临时文件后 -// 重命名为目标文件名,再交给 115 上传(/open/upload/init 的 file_name 取真实文件名)。 +// 为防止多并发上传线程在同一临时目录下发生同名文件(如 poster.jpg)碰撞覆盖与误删, +// 为每个上传任务分配专属临时子目录。 func (p *openAPI115Provider) PutFileNamed(ctx context.Context, parentCID, fileName string, r io.Reader) error { - tmp, err := os.CreateTemp("", "mebox-upload-*") + tmpDir, err := os.MkdirTemp("", "mebox-upload-*") + if err != nil { + return fmt.Errorf("115: 创建临时目录失败:%w", err) + } + defer func() { + _ = os.RemoveAll(tmpDir) + }() + + safeName := filepath.Base(fileName) + if safeName == "" || safeName == "." { + safeName = "file" + } + tmpPath := filepath.Join(tmpDir, safeName) + dst, err := os.OpenFile(tmpPath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644) if err != nil { return fmt.Errorf("115: 创建临时文件失败:%w", err) } - tmpPath := tmp.Name() - defer func() { - _ = tmp.Close() - _ = os.Remove(tmpPath) - }() - if _, err := io.Copy(tmp, r); err != nil { + if _, err := io.Copy(dst, r); err != nil { + _ = dst.Close() return fmt.Errorf("115: 写入临时文件失败:%w", err) } - if err := tmp.Close(); err != nil { + if err := dst.Close(); err != nil { return fmt.Errorf("115: 关闭临时文件失败:%w", err) } - // 重命名为目标文件名,保证上传到 115 后保留原始文件名。 - // 重命名失败必须 fail fast:静默用随机临时名上传会导致 115 上的文件名 - // 变成 mebox-upload-xxx,破坏元数据文件名契约。 - if fileName != "" && fileName != filepath.Base(tmpPath) { - namedPath := filepath.Join(filepath.Dir(tmpPath), fileName) - if err := os.Rename(tmpPath, namedPath); err != nil { - return fmt.Errorf("115: 重命名临时文件为 %s 失败:%w", fileName, err) - } - tmpPath = namedPath - } + _, err = p.c.Upload(ctx, tmpPath, parentCID, "", "") if err != nil { return err diff --git a/internal/service/strm_queue.go b/internal/service/strm_queue.go index c56004d..3b63c52 100644 --- a/internal/service/strm_queue.go +++ b/internal/service/strm_queue.go @@ -378,27 +378,39 @@ func (s *StrmService) processUpload115(ctx context.Context, task *model.StrmUplo return } refs := strings.Split(task.RemoteRef, ",") - if err := open115.OpenClient().DeleteFiles(ctx, task.RemotePath, refs...); err != nil { - s.log.Warn("删除网盘旧元数据失败,跳过删除继续上传新文件", - zap.String("task_id", task.ID), - zap.String("local_path", task.LocalPath), - zap.Error(err)) - // 不 return:继续上传新文件,旧副本由下次同步清理 + if err := open115.OpenClient().DeleteFiles(ctx, task.RemotePath, refs...); err != nil { + s.log.Warn("删除网盘旧元数据失败,跳过删除继续上传新文件", + zap.String("task_id", task.ID), + zap.String("local_path", task.LocalPath), + zap.Error(err)) + // 不 return:继续上传新文件,旧副本由下次同步清理 + } + } + // 优先使用直接本地文件上传接口,零拷贝且彻底根除并发临时文件同名碰撞 + if localUploader, ok := provider.(interface { + PutLocalFile(ctx context.Context, parentCID, localPath string) error + }); ok { + if err := localUploader.PutLocalFile(ctx, task.RemotePath, task.LocalPath); err != nil { + s.uploadTaskFailWithRetry(task, "上传失败:"+err.Error()) + return + } + finish(model.StrmTaskDone, "") + return + } + + f, err := os.Open(task.LocalPath) + if err != nil { + s.uploadTaskFailWithRetry(task, "打开本地文件失败:"+err.Error()) + return + } + if err := named.PutFileNamed(ctx, task.RemotePath, task.FileName, f); err != nil { + _ = f.Close() + s.uploadTaskFailWithRetry(task, "上传失败:"+err.Error()) + return } - } - f, err := os.Open(task.LocalPath) - if err != nil { - s.uploadTaskFailWithRetry(task, "打开本地文件失败:"+err.Error()) - return - } - if err := named.PutFileNamed(ctx, task.RemotePath, task.FileName, f); err != nil { _ = f.Close() - s.uploadTaskFailWithRetry(task, "上传失败:"+err.Error()) - return + finish(model.StrmTaskDone, "") } - _ = f.Close() - finish(model.StrmTaskDone, "") -} // downloadTaskFailWithRetry 下载失败任务按退避重试,超过上限标记 failed。 func (s *StrmService) downloadTaskFailWithRetry(task *model.StrmDownloadTask, message string) {