Compare commits

..

1 Commits

2 changed files with 56 additions and 38 deletions
+26 -20
View File
@@ -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
+30 -18
View File
@@ -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) {