From b810beeff51da7d69ed75d026aec8a6897c3bdca Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Sat, 19 Sep 2026 00:50:22 +0200 Subject: [PATCH] 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 Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/filemanager/workflows/remote_download.go | 67 +++++++++++++++++++ .../workflows/remote_download_test.go | 19 ++++++ 2 files changed, 86 insertions(+) diff --git a/pkg/filemanager/workflows/remote_download.go b/pkg/filemanager/workflows/remote_download.go index 29c5071d..467fd984 100644 --- a/pkg/filemanager/workflows/remote_download.go +++ b/pkg/filemanager/workflows/remote_download.go @@ -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 } diff --git a/pkg/filemanager/workflows/remote_download_test.go b/pkg/filemanager/workflows/remote_download_test.go index 30fa4744..90a814ef 100644 --- a/pkg/filemanager/workflows/remote_download_test.go +++ b/pkg/filemanager/workflows/remote_download_test.go @@ -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{