From 772eb1a56b4d2f135498ee50113ccd38ab4cdf1e Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Sat, 19 Sep 2026 02:44:38 +0200 Subject: [PATCH] feat: prefer the file-hosting node for archive tasks (#39) Compress/extract tasks now seed the node pool's preferred slot with the node hosting the source file's storage policy, so data-heavy work runs where the blob already lives instead of pulling it across nodes. Falls back to weighted selection when the policy is node-less or the node lacks the capability. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/filemanager/workflows/archive.go | 5 ++++ pkg/filemanager/workflows/extract.go | 3 ++ pkg/filemanager/workflows/worfklows.go | 39 ++++++++++++++++++++++++++ 3 files changed, 47 insertions(+) diff --git a/pkg/filemanager/workflows/archive.go b/pkg/filemanager/workflows/archive.go index 53fe21b1..2bef70b5 100644 --- a/pkg/filemanager/workflows/archive.go +++ b/pkg/filemanager/workflows/archive.go @@ -125,6 +125,11 @@ func (m *CreateArchiveTask) Do(ctx context.Context) (task.Status, error) { } m.state = state + // Prefer the node hosting the first source file's storage policy. + if len(m.state.Uris) > 0 { + preferPolicyNode(ctx, dep, &m.state.NodeState, m.state.Uris[0]) + } + // select node node, err := allocateNode(ctx, dep, &m.state.NodeState, types.NodeCapabilityCreateArchive) if err != nil { diff --git a/pkg/filemanager/workflows/extract.go b/pkg/filemanager/workflows/extract.go index 3c379cdf..b145a0aa 100644 --- a/pkg/filemanager/workflows/extract.go +++ b/pkg/filemanager/workflows/extract.go @@ -126,6 +126,9 @@ func (m *ExtractArchiveTask) Do(ctx context.Context) (task.Status, error) { } m.state = state + // Prefer the node hosting the source file's storage policy. + preferPolicyNode(ctx, dep, &m.state.NodeState, m.state.Uri) + // select node node, err := allocateNode(ctx, dep, &m.state.NodeState, types.NodeCapabilityExtractArchive) if err != nil { diff --git a/pkg/filemanager/workflows/worfklows.go b/pkg/filemanager/workflows/worfklows.go index 475dd9be..1356d099 100644 --- a/pkg/filemanager/workflows/worfklows.go +++ b/pkg/filemanager/workflows/worfklows.go @@ -9,7 +9,11 @@ import ( "github.com/cloudreve/Cloudreve/v4/application/dependency" "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/cloudreve/Cloudreve/v4/inventory" "github.com/cloudreve/Cloudreve/v4/pkg/cluster" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/fs/dbfs" + "github.com/cloudreve/Cloudreve/v4/pkg/filemanager/manager" "github.com/cloudreve/Cloudreve/v4/pkg/queue" "github.com/cloudreve/Cloudreve/v4/pkg/util" ) @@ -60,3 +64,38 @@ func prepareTempFolder(ctx context.Context, dep dependency.Dep, t queue.Task) (s dep.Logger().Info("Temp folder created: %s", tempPath) return tempPath, nil } + +// preferPolicyNode seeds the preferred node with the node hosting the given +// file's storage policy, so data-heavy tasks run where the blob already +// lives. No-op when the URI is invalid, has no entity, or the policy is not +// bound to a node — the pool then falls back to weighted selection. +func preferPolicyNode(ctx context.Context, dep dependency.Dep, state *NodeState, uriStr string) { + if state.NodeID > 0 { + return + } + + uri, err := fs.NewUriFromString(uriStr) + if err != nil { + return + } + + fm := manager.NewFileManager(dep, inventory.UserFromContext(ctx)) + defer fm.Recycle() + + file, err := fm.Get(ctx, uri, dbfs.WithFileEntities()) + if err != nil || file == nil { + return + } + + entity := file.PrimaryEntity() + if entity == nil { + return + } + + policy, err := dep.StoragePolicyClient().GetPolicyByID(ctx, entity.PolicyID()) + if err != nil || policy == nil { + return + } + + state.NodeID = policy.NodeID +}