Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
6d4c938
feat(task): add task hook,batch task
sevxn007 Jul 15, 2025
16ed788
Update internal/task/batch_task/refresh.go
sevxn007 Jul 15, 2025
ae92148
fix: upload task allFinish judge
sevxn007 Jul 15, 2025
bbc82a8
Update internal/task/batch_task/refresh.go
sevxn007 Jul 15, 2025
4155767
feat: enhance concurrency safety
sevxn007 Jul 15, 2025
02565a0
优化代码
j2rong4cn Jul 15, 2025
5040c2c
解压缩
j2rong4cn Jul 15, 2025
d85e12e
修复死锁
j2rong4cn Jul 15, 2025
565b540
refactor(move): move as task
sevxn007 Jul 15, 2025
022decb
重构,优化
j2rong4cn Jul 15, 2025
224c38f
.
j2rong4cn Jul 15, 2025
3ecfe6c
优化,修复bug
j2rong4cn Jul 15, 2025
3399cd9
.
j2rong4cn Jul 15, 2025
2be3d8b
修复bug
j2rong4cn Jul 16, 2025
5c3741f
feat: add task retry judge
sevxn007 Jul 17, 2025
f0b4276
代理Task.SetState函数来判断Task的生命周期
j2rong4cn Jul 17, 2025
723045e
chore: use OnSucceeded、OnFailed、OnBeforeRetry functions
sevxn007 Jul 18, 2025
f784784
优化
j2rong4cn Jul 18, 2025
28dc3c0
优化,去除重复代码
j2rong4cn Jul 19, 2025
9601140
Merge remote-tracking branch 'origin/main' into feat/task
j2rong4cn Jul 19, 2025
52f75ee
.
j2rong4cn Jul 19, 2025
9d8470f
优化
j2rong4cn Jul 20, 2025
e9bf00e
Merge remote-tracking branch 'origin/main' into feat/task
j2rong4cn Jul 20, 2025
1ce512c
Merge branch 'OpenListTeam:main' into feat/task
sevxn007 Jul 21, 2025
330fae5
.
j2rong4cn Jul 21, 2025
2de3339
webdav
j2rong4cn Jul 22, 2025
85c2b9c
Revert "fix(fs):After the file is copied or moved, flush the cache of…
j2rong4cn Jul 24, 2025
a4f3a52
Merge branch 'main' into feat/task
j2rong4cn Jul 24, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion drivers/alias/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -193,7 +193,8 @@ func (d *Alias) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
}
if len(srcPath) == len(dstPath) {
for i := range srcPath {
err = errors.Join(err, fs.Move(ctx, *srcPath[i], *dstPath[i]))
_, e := fs.Move(ctx, *srcPath[i], *dstPath[i])
err = errors.Join(err, e)
}
return err
} else {
Expand Down
4 changes: 2 additions & 2 deletions drivers/doubao/util.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import (
"math/rand"
"net/http"
"net/url"
"path/filepath"
stdpath "path"
"sort"
"strconv"
"strings"
Expand Down Expand Up @@ -353,7 +353,7 @@ func (d *Doubao) getUploadConfig(upConfig *UploadConfig, dataType string, file m
"ServiceId": d.UploadToken.Alice[dataType].ServiceID,
"NeedFallback": "true",
"FileSize": strconv.FormatInt(file.GetSize(), 10),
"FileExtension": filepath.Ext(file.GetName()),
"FileExtension": stdpath.Ext(file.GetName()),
"s": randomString(),
}
}
Expand Down
4 changes: 2 additions & 2 deletions internal/bootstrap/task.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,11 +22,11 @@ func InitTaskManager() {
op.RegisterSettingChangingCallback(func() {
fs.UploadTaskManager.SetWorkersNumActive(taskFilterNegative(setting.GetInt(conf.TaskUploadThreadsNum, conf.Conf.Tasks.Upload.Workers)))
})
fs.CopyTaskManager = tache.NewManager[*fs.CopyTask](tache.WithWorks(setting.GetInt(conf.TaskCopyThreadsNum, conf.Conf.Tasks.Copy.Workers)), tache.WithPersistFunction(db.GetTaskDataFunc("copy", conf.Conf.Tasks.Copy.TaskPersistant), db.UpdateTaskDataFunc("copy", conf.Conf.Tasks.Copy.TaskPersistant)), tache.WithMaxRetry(conf.Conf.Tasks.Copy.MaxRetry))
fs.CopyTaskManager = tache.NewManager[*fs.FileTransferTask](tache.WithWorks(setting.GetInt(conf.TaskCopyThreadsNum, conf.Conf.Tasks.Copy.Workers)), tache.WithPersistFunction(db.GetTaskDataFunc("copy", conf.Conf.Tasks.Copy.TaskPersistant), db.UpdateTaskDataFunc("copy", conf.Conf.Tasks.Copy.TaskPersistant)), tache.WithMaxRetry(conf.Conf.Tasks.Copy.MaxRetry))
op.RegisterSettingChangingCallback(func() {
fs.CopyTaskManager.SetWorkersNumActive(taskFilterNegative(setting.GetInt(conf.TaskCopyThreadsNum, conf.Conf.Tasks.Copy.Workers)))
})
fs.MoveTaskManager = tache.NewManager[*fs.MoveTask](tache.WithWorks(setting.GetInt(conf.TaskMoveThreadsNum, conf.Conf.Tasks.Move.Workers)), tache.WithPersistFunction(db.GetTaskDataFunc("move", conf.Conf.Tasks.Move.TaskPersistant), db.UpdateTaskDataFunc("move", conf.Conf.Tasks.Move.TaskPersistant)), tache.WithMaxRetry(conf.Conf.Tasks.Move.MaxRetry))
fs.MoveTaskManager = tache.NewManager[*fs.FileTransferTask](tache.WithWorks(setting.GetInt(conf.TaskMoveThreadsNum, conf.Conf.Tasks.Move.Workers)), tache.WithPersistFunction(db.GetTaskDataFunc("move", conf.Conf.Tasks.Move.TaskPersistant), db.UpdateTaskDataFunc("move", conf.Conf.Tasks.Move.TaskPersistant)), tache.WithMaxRetry(conf.Conf.Tasks.Move.MaxRetry))
op.RegisterSettingChangingCallback(func() {
fs.MoveTaskManager.SetWorkersNumActive(taskFilterNegative(setting.GetInt(conf.TaskMoveThreadsNum, conf.Conf.Tasks.Move.Workers)))
})
Expand Down
144 changes: 87 additions & 57 deletions internal/fs/archive.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,9 @@ import (
"fmt"
"io"
"math/rand"
"mime"
"os"
stdpath "path"
"path/filepath"
"strconv"
"strings"
"time"

Expand All @@ -21,30 +19,22 @@ import (
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/OpenList/v4/internal/stream"
"github.com/OpenListTeam/OpenList/v4/internal/task"
"github.com/OpenListTeam/OpenList/v4/internal/task_group"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/server/common"
"github.com/OpenListTeam/tache"
"github.com/pkg/errors"
log "github.com/sirupsen/logrus"
)

type ArchiveDownloadTask struct {
task.TaskExtension
TaskData
model.ArchiveDecompressArgs
status string
SrcObjPath string
DstDirPath string
srcStorage driver.Driver
dstStorage driver.Driver
SrcStorageMp string
DstStorageMp string
}

func (t *ArchiveDownloadTask) GetName() string {
return fmt.Sprintf("decompress [%s](%s)[%s] to [%s](%s) with password <%s>", t.SrcStorageMp, t.SrcObjPath,
t.InnerPath, t.DstStorageMp, t.DstDirPath, t.Password)
}

func (t *ArchiveDownloadTask) GetStatus() string {
return t.status
return fmt.Sprintf("decompress [%s](%s)[%s] to [%s](%s) with password <%s>", t.SrcStorageMp, t.SrcActualPath,
t.InnerPath, t.DstStorageMp, t.DstActualPath, t.Password)
}

func (t *ArchiveDownloadTask) Run() error {
Expand All @@ -58,16 +48,21 @@ func (t *ArchiveDownloadTask) Run() error {
if err != nil {
return err
}
uploadTask.groupID = stdpath.Join(uploadTask.DstStorageMp, uploadTask.DstActualPath)
task_group.TransferCoordinator.AddTask(uploadTask.groupID, nil)
ArchiveContentUploadTaskManager.Add(uploadTask)
return nil
}

func (t *ArchiveDownloadTask) RunWithoutPushUploadTask() (*ArchiveContentUploadTask, error) {
var err error
if t.srcStorage == nil {
t.srcStorage, err = op.GetStorageByMountPath(t.SrcStorageMp)
if t.SrcStorage == nil {
t.SrcStorage, err = op.GetStorageByMountPath(t.SrcStorageMp)
if err != nil {
return nil, err
}
}
srcObj, tool, ss, err := op.GetArchiveToolAndStream(t.Ctx(), t.srcStorage, t.SrcObjPath, model.LinkArgs{})
srcObj, tool, ss, err := op.GetArchiveToolAndStream(t.Ctx(), t.SrcStorage, t.SrcActualPath, model.LinkArgs{})
if err != nil {
return nil, err
}
Expand All @@ -87,7 +82,7 @@ func (t *ArchiveDownloadTask) RunWithoutPushUploadTask() (*ArchiveContentUploadT
total += s.GetSize()
}
t.SetTotalBytes(total)
t.status = "getting src object"
t.Status = "getting src object"
for _, s := range ss {
if s.GetFile() == nil {
_, err = stream.CacheFullInTempFileAndWriter(s, func(p float64) {
Expand All @@ -104,7 +99,7 @@ func (t *ArchiveDownloadTask) RunWithoutPushUploadTask() (*ArchiveContentUploadT
} else {
decompressUp = t.SetProgress
}
t.status = "walking and decompressing"
t.Status = "walking and decompressing"
dir, err := os.MkdirTemp(conf.Conf.TempDir, "dir-*")
if err != nil {
return nil, err
Expand All @@ -117,13 +112,14 @@ func (t *ArchiveDownloadTask) RunWithoutPushUploadTask() (*ArchiveContentUploadT
uploadTask := &ArchiveContentUploadTask{
TaskExtension: task.TaskExtension{
Creator: t.GetCreator(),
ApiUrl: t.ApiUrl,
},
ObjName: baseName,
InPlace: !t.PutIntoNewDir,
FilePath: dir,
DstDirPath: t.DstDirPath,
dstStorage: t.dstStorage,
DstStorageMp: t.DstStorageMp,
ObjName: baseName,
InPlace: !t.PutIntoNewDir,
FilePath: dir,
DstActualPath: t.DstActualPath,
dstStorage: t.DstStorage,
DstStorageMp: t.DstStorageMp,
}
return uploadTask, nil
}
Expand All @@ -132,18 +128,19 @@ var ArchiveDownloadTaskManager *tache.Manager[*ArchiveDownloadTask]

type ArchiveContentUploadTask struct {
task.TaskExtension
status string
ObjName string
InPlace bool
FilePath string
DstDirPath string
dstStorage driver.Driver
DstStorageMp string
finalized bool
status string
ObjName string
InPlace bool
FilePath string
DstActualPath string
dstStorage driver.Driver
DstStorageMp string
finalized bool
groupID string
}

func (t *ArchiveContentUploadTask) GetName() string {
return fmt.Sprintf("upload %s to [%s](%s)", t.ObjName, t.DstStorageMp, t.DstDirPath)
return fmt.Sprintf("upload %s to [%s](%s)", t.ObjName, t.DstStorageMp, t.DstActualPath)
}

func (t *ArchiveContentUploadTask) GetStatus() string {
Expand All @@ -163,21 +160,42 @@ func (t *ArchiveContentUploadTask) Run() error {
})
}

func (t *ArchiveContentUploadTask) RunWithNextTaskCallback(f func(nextTsk *ArchiveContentUploadTask) error) error {
func (t *ArchiveContentUploadTask) OnSucceeded() {
task_group.TransferCoordinator.Done(t.groupID, true)
}

func (t *ArchiveContentUploadTask) OnFailed() {
task_group.TransferCoordinator.Done(t.groupID, false)
}

func (t *ArchiveContentUploadTask) SetRetry(retry int, maxRetry int) {
t.TaskExtension.SetRetry(retry, maxRetry)
if retry == 0 &&
(len(t.groupID) == 0 || // 重启恢复
(t.GetErr() == nil && t.GetState() != tache.StatePending)) { // 手动重试
t.groupID = stdpath.Join(t.DstStorageMp, t.DstActualPath)
task_group.TransferCoordinator.AddTask(t.groupID, nil)
}
}

func (t *ArchiveContentUploadTask) RunWithNextTaskCallback(f func(nextTask *ArchiveContentUploadTask) error) error {
var err error
if t.dstStorage == nil {
t.dstStorage, err = op.GetStorageByMountPath(t.DstStorageMp)
if err != nil {
return err
}
}
info, err := os.Stat(t.FilePath)
if err != nil {
return err
}
if info.IsDir() {
t.status = "src object is dir, listing objs"
nextDstPath := t.DstDirPath
nextDstActualPath := t.DstActualPath
if !t.InPlace {
nextDstPath = stdpath.Join(nextDstPath, t.ObjName)
err = op.MakeDir(t.Ctx(), t.dstStorage, nextDstPath)
nextDstActualPath = stdpath.Join(nextDstActualPath, t.ObjName)
err = op.MakeDir(t.Ctx(), t.dstStorage, nextDstActualPath)
if err != nil {
return err
}
Expand All @@ -186,6 +204,9 @@ func (t *ArchiveContentUploadTask) RunWithNextTaskCallback(f func(nextTsk *Archi
if err != nil {
return err
}
if !t.InPlace && len(t.groupID) > 0 {
task_group.TransferCoordinator.AppendPayload(t.groupID, task_group.DstPathToRefresh(nextDstActualPath))
}
var es error
for _, entry := range entries {
var nextFilePath string
Expand All @@ -198,16 +219,21 @@ func (t *ArchiveContentUploadTask) RunWithNextTaskCallback(f func(nextTsk *Archi
es = stderrors.Join(es, err)
continue
}
if len(t.groupID) > 0 {
task_group.TransferCoordinator.AddTask(t.groupID, nil)
}
err = f(&ArchiveContentUploadTask{
TaskExtension: task.TaskExtension{
Creator: t.GetCreator(),
ApiUrl: t.ApiUrl,
},
ObjName: entry.Name(),
InPlace: false,
FilePath: nextFilePath,
DstDirPath: nextDstPath,
dstStorage: t.dstStorage,
DstStorageMp: t.DstStorageMp,
ObjName: entry.Name(),
InPlace: false,
FilePath: nextFilePath,
DstActualPath: nextDstActualPath,
dstStorage: t.dstStorage,
DstStorageMp: t.DstStorageMp,
groupID: t.groupID,
})
if err != nil {
es = stderrors.Join(es, err)
Expand All @@ -228,13 +254,13 @@ func (t *ArchiveContentUploadTask) RunWithNextTaskCallback(f func(nextTsk *Archi
Size: info.Size(),
Modified: time.Now(),
},
Mimetype: mime.TypeByExtension(filepath.Ext(t.ObjName)),
Mimetype: utils.GetMimeType(stdpath.Ext(t.ObjName)),
WebPutAsTask: true,
Reader: file,
}
fs.Closers.Add(file)
t.status = "uploading"
err = op.Put(t.Ctx(), t.dstStorage, t.DstDirPath, fs, t.SetProgress, true)
err = op.Put(t.Ctx(), t.dstStorage, t.DstActualPath, fs, t.SetProgress, true)
if err != nil {
return err
}
Expand Down Expand Up @@ -271,8 +297,9 @@ func moveToTempPath(path, prefix string) (string, error) {

func genTempFileName(prefix string) (string, error) {
retry := 0
t := time.Now().UnixMilli()
for retry < 10000 {
newPath := stdpath.Join(conf.Conf.TempDir, prefix+strconv.FormatUint(uint64(rand.Uint32()), 10))
newPath := filepath.Join(conf.Conf.TempDir, prefix+fmt.Sprintf("%x-%x", t, rand.Uint32()))
if _, err := os.Stat(newPath); err != nil {
if os.IsNotExist(err) {
return newPath, nil
Expand Down Expand Up @@ -354,16 +381,19 @@ func archiveDecompress(ctx context.Context, srcObjPath, dstDirPath string, args
}
taskCreator, _ := ctx.Value(conf.UserKey).(*model.User)
tsk := &ArchiveDownloadTask{
TaskExtension: task.TaskExtension{
Creator: taskCreator,
TaskData: TaskData{
TaskExtension: task.TaskExtension{
Creator: taskCreator,
ApiUrl: common.GetApiUrl(ctx),
},
SrcStorage: srcStorage,
DstStorage: dstStorage,
SrcActualPath: srcObjActualPath,
DstActualPath: dstDirActualPath,
SrcStorageMp: srcStorage.GetStorage().MountPath,
DstStorageMp: dstStorage.GetStorage().MountPath,
},
ArchiveDecompressArgs: args,
srcStorage: srcStorage,
dstStorage: dstStorage,
SrcObjPath: srcObjActualPath,
DstDirPath: dstDirActualPath,
SrcStorageMp: srcStorage.GetStorage().MountPath,
DstStorageMp: dstStorage.GetStorage().MountPath,
}
if ctx.Value(conf.NoTaskKey) != nil {
uploadTask, err := tsk.RunWithoutPushUploadTask()
Expand Down
Loading