alist/internal/aria2/monitor.go

186 lines
4.6 KiB
Go
Raw Normal View History

2022-06-22 07:16:13 +00:00
package aria2
import (
"fmt"
"os"
"path"
"path/filepath"
2022-06-22 07:16:13 +00:00
"strconv"
"sync"
"sync/atomic"
"time"
2022-07-31 13:42:01 +00:00
"github.com/alist-org/alist/v3/internal/model"
"github.com/alist-org/alist/v3/internal/op"
2022-07-31 13:42:01 +00:00
"github.com/alist-org/alist/v3/pkg/task"
"github.com/alist-org/alist/v3/pkg/utils"
2022-07-31 13:42:01 +00:00
"github.com/pkg/errors"
log "github.com/sirupsen/logrus"
2022-06-22 07:16:13 +00:00
)
type Monitor struct {
2022-06-22 11:28:41 +00:00
tsk *task.Task[string]
2022-06-22 07:16:13 +00:00
tempDir string
retried int
c chan int
dstDirPath string
2022-06-29 10:12:31 +00:00
finish chan struct{}
2022-06-22 07:16:13 +00:00
}
func (m *Monitor) Loop() error {
defer func() {
notify.Signals.Delete(m.tsk.ID)
// clear temp dir, should do while complete
//_ = os.RemoveAll(m.tempDir)
}()
m.c = make(chan int)
2022-06-29 10:12:31 +00:00
m.finish = make(chan struct{})
2022-06-22 07:16:13 +00:00
notify.Signals.Store(m.tsk.ID, m.c)
2022-06-29 10:12:31 +00:00
var (
err error
ok bool
)
outer:
2022-06-22 07:16:13 +00:00
for {
select {
case <-m.tsk.Ctx.Done():
_, err := client.Remove(m.tsk.ID)
return err
case <-m.c:
2022-06-29 10:12:31 +00:00
ok, err = m.Update()
2022-06-22 07:16:13 +00:00
if ok {
2022-06-29 10:12:31 +00:00
break outer
2022-06-22 07:16:13 +00:00
}
case <-time.After(time.Second * 2):
2022-06-29 10:12:31 +00:00
ok, err = m.Update()
2022-06-22 07:16:13 +00:00
if ok {
2022-06-29 10:12:31 +00:00
break outer
2022-06-22 07:16:13 +00:00
}
}
}
2022-06-29 10:12:31 +00:00
if err != nil {
return err
}
2022-06-29 12:28:02 +00:00
m.tsk.SetStatus("aria2 download completed, transferring")
2022-06-29 10:12:31 +00:00
<-m.finish
2022-06-29 12:28:02 +00:00
m.tsk.SetStatus("completed")
2022-06-29 10:12:31 +00:00
return nil
2022-06-22 07:16:13 +00:00
}
func (m *Monitor) Update() (bool, error) {
info, err := client.TellStatus(m.tsk.ID)
if err != nil {
m.retried++
log.Errorf("failed to get status of %s, retried %d times", m.tsk.ID, m.retried)
return false, nil
2022-06-22 07:16:13 +00:00
}
if m.retried > 5 {
return true, errors.Errorf("failed to get status of %s, retried %d times", m.tsk.ID, m.retried)
}
m.retried = 0
if len(info.FollowedBy) != 0 {
log.Debugf("followen by: %+v", info.FollowedBy)
2022-06-22 07:16:13 +00:00
gid := info.FollowedBy[0]
notify.Signals.Delete(m.tsk.ID)
oldId := m.tsk.ID
2022-06-22 07:16:13 +00:00
m.tsk.ID = gid
DownTaskManager.RawTasks().Delete(oldId)
DownTaskManager.RawTasks().Store(m.tsk.ID, m.tsk)
2022-06-22 07:16:13 +00:00
notify.Signals.Store(gid, m.c)
return false, nil
2022-06-22 07:16:13 +00:00
}
// update download status
total, err := strconv.ParseUint(info.TotalLength, 10, 64)
if err != nil {
total = 0
}
downloaded, err := strconv.ParseUint(info.CompletedLength, 10, 64)
if err != nil {
downloaded = 0
}
2022-06-29 12:28:02 +00:00
progress := float64(downloaded) / float64(total) * 100
m.tsk.SetProgress(int(progress))
2022-06-22 07:16:13 +00:00
switch info.Status {
case "complete":
err := m.Complete()
return true, errors.WithMessage(err, "failed to transfer file")
case "error":
return true, errors.Errorf("failed to download %s, error: %s", m.tsk.ID, info.ErrorMessage)
case "active":
m.tsk.SetStatus("aria2: " + info.Status)
if info.Seeder == "true" {
err := m.Complete()
return true, errors.WithMessage(err, "failed to transfer file")
}
return false, nil
case "waiting", "paused":
2022-06-22 07:16:13 +00:00
m.tsk.SetStatus("aria2: " + info.Status)
return false, nil
case "removed":
return true, errors.Errorf("failed to download %s, removed", m.tsk.ID)
default:
return true, errors.Errorf("failed to download %s, unknown status %s", m.tsk.ID, info.Status)
}
}
2022-07-31 13:42:01 +00:00
var TransferTaskManager = task.NewTaskManager(3, func(k *uint64) {
2022-06-22 07:16:13 +00:00
atomic.AddUint64(k, 1)
})
func (m *Monitor) Complete() error {
// check dstDir again
storage, dstDirActualPath, err := op.GetStorageAndActualPath(m.dstDirPath)
2022-06-22 07:16:13 +00:00
if err != nil {
2022-07-10 06:45:39 +00:00
return errors.WithMessage(err, "failed get storage")
2022-06-22 07:16:13 +00:00
}
// get files
files, err := client.GetFiles(m.tsk.ID)
log.Debugf("files len: %d", len(files))
2022-06-22 07:16:13 +00:00
if err != nil {
return errors.Wrapf(err, "failed to get files of %s", m.tsk.ID)
}
// upload files
var wg sync.WaitGroup
wg.Add(len(files))
go func() {
wg.Wait()
err := os.RemoveAll(m.tempDir)
2022-06-29 10:12:31 +00:00
m.finish <- struct{}{}
2022-06-22 07:16:13 +00:00
if err != nil {
log.Errorf("failed to remove aria2 temp dir: %+v", err.Error())
}
}()
for i, _ := range files {
file := files[i]
2022-07-31 13:42:01 +00:00
TransferTaskManager.Submit(task.WithCancelCtx(&task.Task[uint64]{
Name: fmt.Sprintf("transfer %s to [%s](%s)", file.Path, storage.GetStorage().MountPath, dstDirActualPath),
2022-06-22 11:28:41 +00:00
Func: func(tsk *task.Task[uint64]) error {
2022-06-22 07:16:13 +00:00
defer wg.Done()
2022-06-23 07:57:36 +00:00
size, _ := strconv.ParseInt(file.Length, 10, 64)
mimetype := utils.GetMimeType(file.Path)
2022-06-22 07:16:13 +00:00
f, err := os.Open(file.Path)
if err != nil {
return errors.Wrapf(err, "failed to open file %s", file.Path)
}
2022-07-01 07:04:02 +00:00
stream := &model.FileStream{
2022-09-02 10:24:14 +00:00
Obj: &model.Object{
2022-06-22 07:16:13 +00:00
Name: path.Base(file.Path),
Size: size,
Modified: time.Now(),
IsFolder: false,
},
ReadCloser: f,
Mimetype: mimetype,
2022-06-22 07:16:13 +00:00
}
relDir, err := filepath.Rel(m.tempDir, filepath.Dir(file.Path))
if err != nil {
log.Errorf("find relation directory error: %v", err)
}
newDistDir := filepath.Join(dstDirActualPath, relDir)
return op.Put(tsk.Ctx, storage, newDistDir, stream, tsk.SetProgress)
2022-06-22 07:16:13 +00:00
},
}))
}
return nil
}