feat(workflow): early per-file transfer for multi-file downloads (#3413)

While a remote download is still running, selected files that reached
100% are now transferred to the user's storage in early batches instead
of waiting for the whole torrent. After an early batch, the task returns
to monitoring; the final pass transfers the remainder. Early-batch
failures resume monitoring rather than failing the task.

Authored By: TDvorak <info@tdvorak.dev>

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
pull/3582/head
Tomas Dvorak 2 weeks ago
parent 791764a66f
commit b810beeff5

@ -60,6 +60,10 @@ type (
GetTaskStatusTried int `json:"get_task_status_tried,omitempty"`
Transferred map[int]interface{} `json:"transferred,omitempty"`
Failed int `json:"failed,omitempty"`
// ResumeMonitorAfterTransfer marks an early transfer batch taken while
// the download is still running: after it completes, return to Monitor
// instead of awaiting seeding.
ResumeMonitorAfterTransfer bool `json:"resume_monitor_after_transfer,omitempty"`
}
)
@ -389,6 +393,13 @@ func (m *RemoteDownloadTask) monitor(ctx context.Context, dep dependency.Dep) (t
m.l.Info("Download task seeding completed")
return task.StatusCompleted, nil
case downloader.StatusDownloading:
if m.hasEarlyTransferCandidates(status) {
m.l.Info("Some files already completed, starting early transfer.")
m.state.Phase = RemoteDownloadTaskPhaseTransfer
m.state.ResumeMonitorAfterTransfer = true
m.ResumeAfter(0)
return task.StatusSuspending, nil
}
m.ResumeAfter(resumeAfter)
return task.StatusSuspending, nil
case downloader.StatusUnknown, downloader.StatusError:
@ -399,6 +410,19 @@ func (m *RemoteDownloadTask) monitor(ctx context.Context, dep dependency.Dep) (t
return task.StatusSuspending, nil
}
// hasEarlyTransferCandidates reports whether any selected file finished
// downloading but has not been transferred to the user's storage yet.
func (m *RemoteDownloadTask) hasEarlyTransferCandidates(status *downloader.TaskStatus) bool {
for _, f := range status.Files {
if f.Selected && f.Progress >= 1 {
if _, ok := m.state.Transferred[f.Index]; !ok {
return true
}
}
}
return false
}
func (m *RemoteDownloadTask) slaveTransfer(ctx context.Context, dep dependency.Dep) (task.Status, error) {
u := inventory.UserFromContext(ctx)
if m.state.Transferred == nil {
@ -429,6 +453,11 @@ func (m *RemoteDownloadTask) slaveTransfer(ctx context.Context, dep dependency.D
continue
}
// During early transfer batches only completed files are picked.
if m.state.ResumeMonitorAfterTransfer && f.Progress < 1 {
continue
}
dst := dstUri.JoinRaw(sanitizeFileName(f.Name))
src := path.Join(m.state.Status.SavePath, f.Name)
payload.Files = append(payload.Files, SlaveUploadEntity{
@ -480,12 +509,30 @@ func (m *RemoteDownloadTask) slaveTransfer(ctx context.Context, dep dependency.D
}
m.l.Warning("Slave task %d failed to transfer %d files, retrying...", slaveTaskId, len(m.state.SlaveUploadState.Files)-len(m.state.SlaveUploadState.Transferred))
if m.state.ResumeMonitorAfterTransfer {
// Early transfer batch failed while the download continues -
// return to monitoring; the final transfer pass retries them.
m.state.ResumeMonitorAfterTransfer = false
m.state.Phase = RemoteDownloadTaskPhaseMonitor
m.ResumeAfter(0)
return task.StatusSuspending, nil
}
return task.StatusError, fmt.Errorf(
"slave task failed to transfer %d files, first 5 errors: %s",
len(m.state.SlaveUploadState.Files)-len(m.state.SlaveUploadState.Transferred),
m.state.SlaveUploadState.First5TransferErrors,
)
} else {
if m.state.ResumeMonitorAfterTransfer {
for i := range m.state.SlaveUploadState.Transferred {
m.state.Transferred[m.state.SlaveUploadState.Files[i].Index] = struct{}{}
}
m.state.SlaveUploadTaskID = 0
m.state.ResumeMonitorAfterTransfer = false
m.state.Phase = RemoteDownloadTaskPhaseMonitor
m.ResumeAfter(0)
return task.StatusSuspending, nil
}
m.state.Phase = RemoteDownloadTaskPhaseAwaitSeeding
m.ResumeAfter(0)
return task.StatusSuspending, nil
@ -519,6 +566,11 @@ func (m *RemoteDownloadTask) masterTransfer(ctx context.Context, dep dependency.
allFiles := make([]downloader.TaskFile, 0, len(m.state.Status.Files))
for _, f := range m.state.Status.Files {
if f.Selected {
// Early transfer batches only pick completed files; files still
// downloading are left for the final transfer pass.
if m.state.ResumeMonitorAfterTransfer && f.Progress < 1 {
continue
}
allFiles = append(allFiles, f)
totalSize += f.Size
totalCount++
@ -622,12 +674,27 @@ func (m *RemoteDownloadTask) masterTransfer(ctx context.Context, dep dependency.
wg.Wait()
if failed > 0 {
if m.state.ResumeMonitorAfterTransfer {
// Early batch failed while the download continues - return to
// monitoring; the final transfer pass retries the failed files.
m.l.Warning("Early transfer batch failed for %d file(s), will retry after download completes.", failed)
m.state.ResumeMonitorAfterTransfer = false
m.state.Phase = RemoteDownloadTaskPhaseMonitor
m.ResumeAfter(0)
return task.StatusSuspending, nil
}
m.state.Failed = int(failed)
m.l.Error("Failed to transfer %d file(s).", failed)
return task.StatusError, fmt.Errorf("failed to transfer %d file(s), first 5 errors: %s", failed, ae.FormatFirstN(5))
}
m.l.Info("All files transferred.")
if m.state.ResumeMonitorAfterTransfer {
m.state.ResumeMonitorAfterTransfer = false
m.state.Phase = RemoteDownloadTaskPhaseMonitor
m.ResumeAfter(0)
return task.StatusSuspending, nil
}
m.state.Phase = RemoteDownloadTaskPhaseAwaitSeeding
return task.StatusSuspending, nil
}

@ -7,6 +7,7 @@ import (
"github.com/cloudreve/Cloudreve/v4/inventory/types"
"github.com/cloudreve/Cloudreve/v4/pkg/cluster"
"github.com/cloudreve/Cloudreve/v4/pkg/downloader"
"github.com/stretchr/testify/assert"
)
@ -129,6 +130,24 @@ func TestBuildDownloadOptionsHeaders(t *testing.T) {
a.False(hasHeader)
}
func TestHasEarlyTransferCandidates(t *testing.T) {
a := assert.New(t)
m := &RemoteDownloadTask{state: &RemoteDownloadTaskState{}}
status := &downloader.TaskStatus{Files: []downloader.TaskFile{
{Index: 0, Selected: true, Progress: 1},
{Index: 1, Selected: true, Progress: 0.5},
{Index: 2, Selected: false, Progress: 1},
}}
a.True(m.hasEarlyTransferCandidates(status))
m.state.Transferred = map[int]interface{}{0: struct{}{}}
a.False(m.hasEarlyTransferCandidates(status))
status.Files[1].Progress = 1
a.True(m.hasEarlyTransferCandidates(status))
}
func TestNewRemoteDownloadTaskSanitizesFileName(t *testing.T) {
a := assert.New(t)
tsk, err := NewRemoteDownloadTask(context.Background(), "https://example.com/f", "", "cloudreve://my/dst", &RemoteDownloadTaskOption{

Loading…
Cancel
Save